为何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()被重复调用。这是它的设计限制,不能代表生产环境的真实表现。
解决建议
- 生产环境验证:如果要确认
setup()的正确行为,直接用Dataflow Runner提交到GCP运行,此时setup()会按预期仅在Worker启动时调用一次,是创建数据库连接的理想位置。 - 本地调试优化:如果需要在DirectRunner中模拟连接复用,可以把数据库连接作为DoFn的实例变量,在
start_bundle()中检查连接状态(比如是否断开),仅在需要时重新创建,避免每次Bundle都新建连接。
内容的提问来源于stack exchange,提问作者Machiel Treffers
相关产品推荐
相关产品推荐

