You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用GCP Dataflow与Cloud Functions处理CSV文件遇模块缺失错误

问题:Dataflow运行时出现ModuleNotFoundError: No module named 'functions_framework'

代码示例(main.py)

import functions_framework
import re
import apache_beam as beam
from apache_beam.options.pipeline_options import PipelineOptions

def process_csv_file(element):
    # Cleans up CSV fields by removing leading/trailing whitespaces and filtering out 'NaN'
    cleaned_data = [field.strip() for field in element.split(',') if field.strip() != 'NaN']
    return ','.join(cleaned_data)

def generate_valid_job_name(file_name):
    # Generates a valid job name by replacing non-alphanumeric characters with underscores
    cleaned_name = re.sub(r'[^a-z0-9]', '_', file_name.lower())
    if not cleaned_name[0].isalpha():
        cleaned_name = 'a' + cleaned_name[1:]
    if not cleaned_name[-1].isalnum():
        cleaned_name = cleaned_name[:-1] + 'z'
    if not cleaned_name:
        cleaned_name = 'default_job_name'
    cleaned_name = cleaned_name.replace('_', '')
    return cleaned_name

@functions_framework.cloud_event
def start_dataflow_process(cloud_event):
    # Imports necessary modules
    import base64
    import json

    # Define pipeline options and parameters
    project_id = 'your_project_id'
    input_bucket = 'your_input_bucket'
    output_bucket = 'your_output_bucket'
    output_prefix = 'Out_'
    
    # Extracts relevant information from the Cloud Event data
    data = cloud_event.data
    file_name = data['name']
    bucket_name = data['bucket']
    
    # Checks if the keyword 'input' is present in the file name
    if 'input' not in file_name.lower():
        print(f"File {file_name} does not contain the keyword 'input'. Exiting.")
        return
    
    print(f"File: {file_name}, Bucket: {bucket_name}")
    
    # Generates a valid job name based on the input file name
    job_name = generate_valid_job_name(file_name)
    
    # Configures Dataflow pipeline options
    pipeline_options = {
        'project': project_id,
        'runner': 'DataflowRunner',
        'staging_location': f'gs://{output_bucket}/{job_name}/staging',
        'temp_location': f'gs://{output_bucket}/{job_name}/temp',
        'job_name': job_name,
        'region': 'your_region',
        'save_main_session': True
    }
    
    # Creates a Dataflow pipeline
    pipeline = beam.Pipeline(options=PipelineOptions.from_dictionary(pipeline_options))
    
    # Reads CSV data from the input file
    csv_data = pipeline | 'ReadFromText' >> beam.io.ReadFromText(f'gs://{input_bucket}/{file_name}')
    
    # Processes CSV data by cleaning up fields
    cleaned_data = csv_data | 'ProcessCSV' >> beam.Map(process_csv_file)
    
    # Defines the output path for the cleaned data
    output_path = f'gs://{bucket_name}/{output_prefix}{file_name}'
    
    # Writes the cleaned data to a text file
    cleaned_data | 'WriteToText' >> beam.io.WriteToText(output_path)
    
    print("Pipeline starting...")
    
    # Runs the Dataflow pipeline
    result = pipeline.run()
    
    print("Pipeline started")
    
    # Waits for the pipeline to finish
    result.wait_until_finish()
    
    print("Pipeline ended")

错误信息

File "/usr/local/lib/python3.8/site-packages/dill/_dill.py", line 827, in _import_module
return getattr(import(module, None, None, [obj]), obj)
ModuleNotFoundError: No module named 'functions_framework'

原因分析

  • functions_framework是GCP Cloud Functions的专属依赖包,Dataflow的Worker节点环境默认不会安装这个包。
  • 代码将Cloud Functions的触发逻辑(@functions_framework.cloud_event装饰的函数)和Dataflow管道逻辑写在同一个文件中,开启save_main_session=True后,Dataflow会序列化整个主会话,包括functions_framework的导入语句,导致Worker节点在反序列化时尝试导入不存在的包,抛出错误。

解决方案

1. 拆分代码结构

将Cloud Functions触发逻辑和Dataflow管道逻辑分离成两个独立文件:

(1)dataflow_pipeline.py(独立的Dataflow管道代码)

import re
import apache_beam as beam
from apache_beam.options.pipeline_options import PipelineOptions

def process_csv_file(element):
    cleaned_data = [field.strip() for field in element.split(',') if field.strip() != 'NaN']
    return ','.join(cleaned_data)

def generate_valid_job_name(file_name):
    cleaned_name = re.sub(r'[^a-z0-9]', '_', file_name.lower())
    if not cleaned_name[0].isalpha():
        cleaned_name = 'a' + cleaned_name[1:]
    if not cleaned_name[-1].isalnum():
        cleaned_name = cleaned_name[:-1] + 'z'
    if not cleaned_name:
        cleaned_name = 'default_job_name'
    cleaned_name = cleaned_name.replace('_', '')
    return cleaned_name

def run_pipeline(project_id, input_bucket, output_bucket, file_name, bucket_name):
    job_name = generate_valid_job_name(file_name)
    pipeline_options = {
        'project': project_id,
        'runner': 'DataflowRunner',
        'staging_location': f'gs://{output_bucket}/{job_name}/staging',
        'temp_location': f'gs://{output_bucket}/{job_name}/temp',
        'job_name': job_name,
        'region': 'your_region',
        'save_main_session': True
    }
    
    pipeline = beam.Pipeline(options=PipelineOptions.from_dictionary(pipeline_options))
    
    csv_data = pipeline | 'ReadFromText' >> beam.io.ReadFromText(f'gs://{input_bucket}/{file_name}')
    cleaned_data = csv_data | 'ProcessCSV' >> beam.Map(process_csv_file)
    
    output_path = f'gs://{bucket_name}/Out_{file_name}'
    cleaned_data | 'WriteToText' >> beam.io.WriteToText(output_path)
    
    result = pipeline.run()
    result.wait_until_finish()

(2)main.py(Cloud Functions触发代码)

import functions_framework
import json
from dataflow_pipeline import run_pipeline

@functions_framework.cloud_event
def start_dataflow_process(cloud_event):
    project_id = 'your_project_id'
    input_bucket = 'your_input_bucket'
    output_bucket = 'your_output_bucket'
    
    data = cloud_event.data
    file_name = data['name']
    bucket_name = data['bucket']
    
    if 'input' not in file_name.lower():
        print(f"File {file_name} does not contain the keyword 'input'. Exiting.")
        return
    
    print(f"File: {file_name}, Bucket: {bucket_name}")
    print("Pipeline starting...")
    
    run_pipeline(project_id, input_bucket, output_bucket, file_name, bucket_name)
    
    print("Pipeline ended")

2. 配置Dataflow依赖(如需额外包)

如果Dataflow管道依赖第三方包,可创建requirements.txt文件:

apache-beam[gcp]==2.54.0

然后在PipelineOptions中添加参数指定依赖文件:

pipeline_options = {
    # ...其他参数
    'requirements_file': './requirements.txt'
}

3. 部署注意事项

  • 部署Cloud Functions时,确保包含dataflow_pipeline.py和requirements.txt(如有)。
  • 确保Dataflow服务账号拥有访问GCS存储桶和Dataflow资源的足够权限。

内容的提问来源于stack exchange,提问作者nikhil9933

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.07 14:14:54