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

Google Dataflow运行报错:write_to_pubsub未定义求助

问题解决:DataflowRunner运行时出现NameError: name 'write_to_pubsub' is not defined

问题概述

你的代码在使用Python DirectRunner时可正常从GCS读取CSV并发布到Pub/Sub,但切换到Google Cloud DataflowRunner运行时,抛出NameError: name 'write_to_pubsub' is not defined错误。

错误原因

  1. Lambda函数序列化缺陷:Dataflow会将代码分发到远程worker节点执行,你使用的lambda rows: write_to_pubsub(rows, publisher)在序列化时无法携带write_to_pubsub函数的引用,导致worker节点找不到该函数。
  2. Pub/Sub客户端无法跨进程序列化:在主进程初始化的pubsub_v1.PublisherClient不能直接传递给worker,客户端对象不支持序列化,会触发运行时错误。
  3. 冗余行拆分逻辑:ReadFromText已实现逐行读取文件内容,后续的FlatMap(lambda content: content.split("\n"))属于重复操作,还可能引入空行问题。

修复方案

  1. 自定义DoFn类:将Pub/Sub发布逻辑封装到继承自beam.DoFn的类中,确保逻辑能被正确序列化到worker。
  2. 在worker进程初始化客户端:在DoFn的setup方法中创建PublisherClient,每个worker进程独立初始化客户端,规避序列化问题。
  3. 移除冗余行拆分步骤:删除Split CSV lines步骤,直接使用ReadFromText的逐行输出。
  4. 明确传递主题参数:将Pub/Sub主题作为参数传入DoFn,避免硬编码或跨进程引用问题。

修复后完整代码

import apache_beam as beam
from apache_beam.options.pipeline_options import PipelineOptions
from google.cloud import pubsub_v1
import json


def read_csv_file(pipeline, file_path):
    """Reads a CSV file and returns a PCollection of rows.

    Args:
        pipeline: A Beam pipeline.
        file_path: The path to the CSV file.

    Returns:
        A PCollection of rows.
    """
    rows = (
        pipeline
        | "Read CSV" >> beam.io.ReadFromText(file_path, skip_header_lines=1)
        | "Parse CSV" >> beam.Map(lambda line: line.split(","))
        | "Create record dictionaries" >> beam.Map(lambda record: {
            'id': int(record[0]),
            'firstname': record[1],
            'lastname': record[2]
         })
    )
    return rows


class WriteToPubSub(beam.DoFn):
    """Custom DoFn to publish records to Pub/Sub."""
    def __init__(self, topic_name):
        self.topic_name = topic_name
        self.publisher = None

    def setup(self):
        # Initialize publisher in worker process setup
        self.publisher = pubsub_v1.PublisherClient()

    def process(self, element):
        data_str = json.dumps(element)
        print(f"Preparing a JSON-encoded message:\n{data_str}")
        data = data_str.encode("utf-8")
        # Publish record to Pub/Sub
        future = self.publisher.publish(self.topic_name, data=data)
        # Optional: Wait for publish to complete (for sync behavior)
        # future.result()


if __name__ == "__main__":
    # Job options
    options = PipelineOptions(
        project='XXX',
        runner='DataflowRunner',
        streaming=True,
        job_name='dataflow-job',
        staging_location='gs://dataflow-bucket/staging',
        temp_location='gs://dataflow-bucket/temp',
        region='europe-west8',
        requirements_file='./requirements.txt'
    )

    # Pub/Sub topic name
    pubsub_topic = "projects/XXX/topics/topic1"

    # Create a pipeline.
    pipeline = beam.Pipeline(options=options)

    # Read the CSV file and publish to pubsub
    rows = read_csv_file(pipeline, "gs://dataflow-bucket/input/*.csv") | beam.ParDo(WriteToPubSub(pubsub_topic))

    # Run the pipeline.
    pipeline.run()

额外注意事项

  • 确保requirements.txt包含必要依赖:apache-beam[gcp]>=2.50.0、google-cloud-pubsub>=2.19.0
  • 若需保证消息发布可靠性,可取消注释future.result(),但会增加处理延迟
  • 生产环境建议添加错误处理逻辑,捕获publish方法可能抛出的异常

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.16 02:14:55