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

为何Apache Beam的DoFn.setup()在Worker启动后多次调用?

关于Apache Beam DoFn.setup()在DirectRunner中重复调用的问题

我在做基于Python的流式Dataflow管道实验,需要把读取的数据流写入CloudSQL的PostgreSQL实例,所以在找创建数据库连接的合适位置。因为用ParDo函数写数据,原本以为DoFn.setup()是合适的选择——查资料说这个方法只会在Worker启动时调用一次。但测试发现,setup()的调用次数和start_bundle()(处理一定数量元素后调用)的次数一样。

我写了个简单管道:从PubSub读取消息,提取对象文件名输出,同时记录setup()和start_bundle()的调用时间(代码如下)。用DirectRunner运行时,预期setup()只记录一次,但实际日志里它的调用次数和start_bundle()完全一致。我想知道这是我对setup()功能理解错了,还是有其他原因?现在看来setup()好像不是创建数据库连接的理想位置。

import argparse
import logging
from datetime import datetime

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

setup_counter=0
bundle_counter=0

class GetFileName(beam.DoFn):
    """
    Generate file path from PubSub message attributes
    """
    def _now(self):
        return datetime.now().strftime("%Y/%m/%d %H:%M:%S")

    def setup(self):
        global setup_counter
                
        moment = self._now()
        logging.info("setup() called %s" % moment)
        
        setup_counter=setup_counter+1
        logging.info(f"setup_counter = {setup_counter}")

    def start_bundle(self):
        global bundle_counter
        
        moment = self._now()
        logging.info("Bundle started %s" % moment)
        
        bundle_counter=bundle_counter+1
        logging.info(f"Bundle_counter = {bundle_counter}")

    def process(self, element):
        attr = dict(element.attributes)

        objectid = attr["objectId"]

        # not sure if this is the prettiest way to create this uri, but works for the poc
        path = f'{objectid}'

        yield path


def run(input_subscription, pipeline_args=None):

    pipeline_options = PipelineOptions(
        pipeline_args, streaming=True
    )

    with beam.Pipeline(options=pipeline_options) as pipeline:

        files = (pipeline
                 | "Read from PubSub" >> beam.io.ReadFromPubSub(subscription=input_subscription,
                                                                with_attributes=True)
                 | "Get filepath" >> beam.ParDo(GetFileName())
                )

        files | "Print results" >> beam.Map(logging.info)


if __name__ == "__main__":
    logging.getLogger().setLevel(logging.INFO)

    parser = argparse.ArgumentParser()
    parser.add_argument(
        "--input_subscription",
        dest="input_subscription",
        required=True,
        help="The Cloud Pub/Sub subscription to read from."
    )

    known_args, pipeline_args = parser.parse_known_args()

    run(
        known_args.input_subscription,
        pipeline_args
    )

运行日志

python main.py \
>     --runner DirectRunner \
>     --input_subscription <my_subscription> \
>     --direct_num_workers 1 \
>     --streaming true
...
INFO:root:setup() called 2022/11/16 15:11:13
INFO:root:setup_counter = 1
INFO:root:Bundle started 2022/11/16 15:11:13
INFO:root:Bundle_counter = 1
INFO:root:avro/20221116135543584-hlgeinp.avro
...
INFO:root:setup() called 2022/11/16 15:11:16
INFO:root:setup_counter = 2
INFO:root:Bundle started 2022/11/16 15:11:16
INFO:root:Bundle_counter = 2
...
INFO:root:setup() called 2022/11/16 15:11:18
INFO:root:setup_counter = 3
INFO:root:Bundle started 2022/11/16 15:11:18
INFO:root:Bundle_counter = 3
...

原因分析与解决方案

你遇到的setup()重复调用是DirectRunner的本地调试特性导致的,和生产环境的Dataflow Runner行为完全不同:

  • 分布式Dataflow Runner中,setup()确实会在每个Worker进程启动时仅调用一次,适合初始化全局资源(比如数据库长连接),同一个Worker上的所有Bundle都会复用这个资源。
  • DirectRunner是为本地调试设计的,它的实现逻辑会为每个Bundle重新创建DoFn实例,导致setup()被重复调用。这是它的设计限制,不能代表生产环境的真实表现。

解决建议

  1. 生产环境验证:如果要确认setup()的正确行为,直接用Dataflow Runner提交到GCP运行,此时setup()会按预期仅在Worker启动时调用一次,是创建数据库连接的理想位置。
  2. 本地调试优化:如果需要在DirectRunner中模拟连接复用,可以把数据库连接作为DoFn的实例变量,在start_bundle()中检查连接状态(比如是否断开),仅在需要时重新创建,避免每次Bundle都新建连接。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.12 21:21:12