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

Dataflow运行报错'read_inputpath_date未定义',Direct Runner正常

问题分析:Dataflow Runner中函数未定义错误

你的代码在Direct Runner运行正常,但切换到Dataflow Runner时出现'name 'read_inputpath_date' is not defined'错误,核心原因是Beam分布式执行的序列化机制与本地单进程执行的差异,具体如下:

错误原因

Direct Runner在本地单进程环境运行,所有函数定义都在同一个内存空间中,process_messages调用read_inputpath_date时可以直接找到函数引用。但Dataflow Runner是分布式架构:

  • 你的Pipeline代码会被序列化后分发到多个Worker节点执行。
  • 当process_messages作为ParDo的处理函数被序列化时,它对read_inputpath_date的引用是以函数名称的形式保存的,而非函数本身的代码。
  • Worker节点反序列化后,若无法在其Python环境中找到read_inputpath_date的定义(比如代码打包不完整、函数作用域不可见),就会抛出名称未定义的错误。

另外,你直接使用google.cloud.storage客户端读取GCS文件的方式,不符合Beam的分布式设计模式,虽然不是直接引发当前错误的原因,但会带来并行性、容错性等潜在问题。

解决方案

方案1:将函数整合到DoFn类中

把read_inputpath_date作为方法整合到自定义的DoFn类中,确保整个处理逻辑是可序列化的,Worker能正确加载所有方法:

from apache_beam.options.pipeline_options import PipelineOptions
import os
import logging
import apache_beam as beam
import gzip
import json, io

from google.cloud import storage

os.environ["GOOGLE_APPLICATION_CREDENTIALS"] = "credentials.json"

class ProcessMessagesDoFn(beam.DoFn):
    def read_inputpath_date(self, input_path):
        # Read the contents of the .jsonl.gz file from GCS
        storage_client = storage.Client()
        bucket_name, file_name = input_path[5:].split("/", 1)
        bucket = storage_client.bucket(bucket_name)
        blob = bucket.blob(file_name)
        file_content = blob.download_as_bytes()
        data = []
        with gzip.GzipFile(fileobj=io.BytesIO(file_content), mode='r') as f:
            for line in f:
                record = json.loads(line)
                data.append(record)
        logging.info('Done reading data from path{}'.format(input_path))
        return data

    def process(self, element):
        # Process the received Pub/Sub message
        message = json.loads(element.decode('utf-8'))
        logging.info('Received message: {}'.format(message))
        input_path = message["file_path"]
        input_path = 'gs://mf-staging-area/' + input_path
        table_name = message["table_name"]
        logging.info('Input path: {}'.format(input_path))
        logging.info('Table name: {}'.format(table_name))
        try:
            data = self.read_inputpath_date(input_path)
            for line in data:
                table_row = line['payload']
                logging.info('Line: {}'.format(line))
                logging.info('Table: {}'.format(table_name))
                logging.info('Row: {}'.format(table_row))
        except:
            logging.exception('Error occurred while reading input path: {}'.format(input_path))

def run():
    project_id = 'moneyfellows-data'
    subscription = 'projects/moneyfellows-data/subscriptions/streaming-topic-sub'

    pipeline_options = PipelineOptions()

    with beam.Pipeline(options=pipeline_options) as p:
        messages = (
                p
                | 'ReadFromPubSub' >> beam.io.ReadFromPubSub(subscription=subscription)
        )

        processed_messages = (
                messages
                | 'ProcessMessages' >> beam.ParDo(ProcessMessagesDoFn())
        )

if __name__ == '__main__':
    logging.basicConfig(level=logging.INFO)
    run()

方案2:使用Beam内置GCS读取API(推荐)

替换google.cloud.storage客户端为Beam原生的GCS读取组件,更适配分布式环境,同时避免手动处理文件下载和解析的繁琐:

from apache_beam.options.pipeline_options import PipelineOptions
import os
import logging
import apache_beam as beam
import json

os.environ["GOOGLE_APPLICATION_CREDENTIALS"] = "credentials.json"

def process_message_and_read_gcs(element):
    message = json.loads(element.decode('utf-8'))
    logging.info('Received message: {}'.format(message))
    input_path = 'gs://mf-staging-area/' + message["file_path"]
    table_name = message["table_name"]
    logging.info('Input path: {}'.format(input_path))
    logging.info('Table name: {}'.format(table_name))
    
    # 使用Beam的ReadFromText读取压缩文件,自动处理gzip
    with beam.Pipeline() as sub_pipeline:
        records = (
            sub_pipeline
            | 'ReadGCSFile' >> beam.io.ReadFromText(input_path, compression_type='gzip')
            | 'ParseJSON' >> beam.Map(json.loads)
        )
        # 收集数据(仅用于示例,实际可直接在子Pipeline中处理)
        data = list(records)
    
    for line in data:
        table_row = line['payload']
        logging.info('Line: {}'.format(line))
        logging.info('Table: {}'.format(table_name))
        logging.info('Row: {}'.format(table_row))

def run():
    project_id = 'moneyfellows-data'
    subscription = 'projects/moneyfellows-data/subscriptions/streaming-topic-sub'

    pipeline_options = PipelineOptions()

    with beam.Pipeline(options=pipeline_options) as p:
        messages = (
                p
                | 'ReadFromPubSub' >> beam.io.ReadFromPubSub(subscription=subscription)
        )

        processed_messages = (
                messages
                | 'ProcessMessages' >> beam.Map(process_message_and_read_gcs)
        )

if __name__ == '__main__':
    logging.basicConfig(level=logging.INFO)
    run()

额外注意事项

  • 提交Dataflow作业时,确保所有依赖包(如google-cloud-storage)被正确安装,可通过requirements.txt文件配合--requirements_file参数指定。
  • 避免在DoFn中创建全局资源(如storage.Client),建议在DoFn的setup方法中初始化,减少资源开销。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 14:15:01