Run apache beam in GCP Errror ModuleNotFoundError: No module named

Viewed 365

My project is about running an etl developed in apache beam that transforms text files. when I run the pipelinestransform.py file from my python virtual environment Pipelines/pipelinestransform.py it turns out that the execution works. but when I call the execution of this etl (module) using the Runnerdataflow configuration from another main.py file it generates this error

ile "/usr/local/lib/python3.8/site-packages/dill/_dill.py", line 462, in find_class return StockUnpickler.find_class (self, module, name) ModuleNotFoundError: No module named 'Pipelines'

this is my project structure

\myproject
  setup.py
  requirements.txt  
  main.py
  \Pipelines
    __init__.py 
    pipelinestransform.py

pipelinestransform.py

from typing import Tuple

import os 
import apache_beam as beam
import argparse
from apache_beam.options.pipeline_options import PipelineOptions
from apache_beam.options.pipeline_options import SetupOptions
from apache_beam.options.pipeline_options import GoogleCloudOptions
from apache_beam.options.pipeline_options import StandardOptions
from apache_beam.io import ReadFromText
from apache_beam.io import WriteToText


def sanitizar_palabra(palabra):
  para_quitar = [',', '.', '-', ':', ' ', "'", '"']

  for simbolo in para_quitar:
    palabra = palabra.replace(simbolo, '')

  palabra = palabra.lower()

  palabra = palabra.replace("á", "a")
  palabra = palabra.replace("é", "e")
  palabra = palabra.replace("í", "i")
  palabra = palabra.replace("ó", "o")
  palabra = palabra.replace("ú", "u")

  return palabra


def main():
  parser = argparse.ArgumentParser(description="Nuestro primer pipeline")
  parser.add_argument("--entrada", help="Fichero de entrada")
  parser.add_argument("--salida", help="Fichero de salida")
  parser.add_argument("--n-palabras", type=int,  help="Número de palabras en la salida")

  our_args, beam_args = parser.parse_known_args()
  run_pipeline()


def run_pipeline():
      

  entrada = 'gs://project-test-001/data/muestra.txt'
  salida = 'gs://project-test-001/out/salida.csv'
  n_palabras = 50


  #os.path.join(os.path.dirname(__file__), "setup.py")

  pipeline_options = PipelineOptions()

  google_cloud_options = pipeline_options.view_as(GoogleCloudOptions)
  google_cloud_options.project = 'My-project-dev'
  google_cloud_options.job_name = 'datapipeline-conteo'
  google_cloud_options.staging_location = 'gs://project-test-001/dataflow-staging' 
  google_cloud_options.region = 'us-east1'
  google_cloud_options.temp_location = 'gs://project-test-001/tmp/' 
  pipeline_options.view_as(StandardOptions).runner = 'DataFlowRunner'



  with beam.Pipeline(options=pipeline_options) as p:
    lineas: PCollection[str] = p | "Leemos entrada" >> beam.io.ReadFromText(entrada)
    # "En un lugar de La Mancha" --> ["En", "un", ...], [...], [...] --> "En", "un", "lugar", ....
    palabras = lineas | "Pasamos a palabras" >> beam.FlatMap(lambda l: l.split())
    limpiadas = palabras | "Sanitizamos" >> beam.Map(sanitizar_palabra)
    contadas: PCollection[Tuple[str, int]] = limpiadas  | "Contamos" >> beam.combiners.Count.PerElement()
    # "En" -> ("En", 17)
    # "un" -> ("un", 28)
    palabras_top_lista = contadas | "Ranking" >> beam.combiners.Top.Of(n_palabras, key=lambda kv: kv[1])
    palabras_top = palabras_top_lista | "Desenvuelve lista" >> beam.FlatMap(lambda x: x)
    formateado: PCollection[str] = palabras_top | "Formateamos" >> beam.Map(lambda kv: "%s,%d" % (kv[0], kv[1]))
    formateado | "Escribimos salida" >> beam.io.WriteToText(salida)

  # Used while testing locally
  #if pipeline_options.view_as(StandardOptions).runner == "DirectRunner":
  #p.wait_until_finish()

if __name__ == '__main__':
  main()

main.py

import sys
from logging import log
import re
import argparse
import json
import time
from decouple import config
#import pipelines
import logging
from Pipelines.pipelinestransform import *

run_pipeline()
0 Answers
Related