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()