DataflowRunner运行Beam管道报错:'PDone' object has no attribute 'windowing'
问题描述
开发了一个Apache Beam管道,从两个Postgres CloudSQL数据库读取记录,经数据转换后通过WriteToPubSub模块推送到Google PubSub。本地使用DirectRunner运行时,CloudSQL连接和PubSub推送均正常,但设置runner='DataflowRunner'后,管道在Beam的ptransform.py模块的get_windowing函数中报错:'PDone' object has no attribute 'windowing'。不确定Dataflow Runner引入的差异,怀疑是权限问题或需指定窗口选项,核心代码片段如下:
# Invoker code if __name__ == '__main__': # Set up your PostgreSQL connection parameters db_config = { 'host': os.getenv('DB_HOST'), 'port': os.getenv('DB_PORT'), 'database': os.getenv('DB_NAME'), 'user': os.getenv('DB_USER'), 'password': os.getenv('DB_PASS') } parser = argparse.ArgumentParser() args, beam_args = parser.parse_known_args() print(args) publish_topic = os.getenv('PUBLISH_TOPIC') # Set up Apache Beam pipeline options pipeline_options = PipelineOptions( beam_args, runner='DataflowRunner', project='<gcp-project-id>', job_name='dispatch-demo-1', temp_location='<bucket path>', region='europe-west1') dispatch_args = pipeline_options.view_as(DispatchOptions) with beam.Pipeline(options=pipeline_options) as pipeline: # Create a dummy input element dummy_input = pipeline | beam.Create(['dummy']) client_operations_query = f"SELECT * FROM table2 WHERE attr1=abc" get_retailer_categories_query = 'SELECT * FROM table1' co_rows = dummy_input | 'Get table1 rows' >> ReadDB( client_operations_query, **db_config) | 'CO list to map' >> beam.Map(list_to_dict, 'internal_category_id') co_rows | "co_rows " >> beam.Map(print) retailer_cat_rows = dummy_input | 'Get table2 rows' >> ReadDB( get_retailer_categories_query, **db_config) | 'Table2 list to map' >> beam.Map(list_to_dict, 'internal_category_id') retailer_cat_rows | "retailer_cat_rows" >> beam.Map(print) denormalised_co_rows = (({ 'co_rows': co_rows, 'retailer_cat_rows': retailer_cat_rows }) | 'group by cat_ids' >> beam.CoGroupByKey() | 'Join by cat_id' >> beam.ParDo(MergeTransform())) groupedRows = denormalised_co_rows | beam.GroupBy(get_hash) | 'ExtractClientIds' >> beam.Map(lambda element: (element[0], [obj['client_id'] for obj in element[1]], element[1])) | "Convert to string" >> beam.Map( encode_as_task) | "Write to Pub/Sub" >> beam.io.WriteToPubSub(topic=publish_topic) | "pubsub out print" >> beam.Map(print)
问题原因与解决方法
这个错误的核心原因是**WriteToPubSub返回的是PDone类型,该类型没有windowing属性,后续无法挂载需要依赖窗口信息的变换(比如beam.Map(print))**。DirectRunner对这种非规范写法容忍度较高,但DataflowRunner的校验逻辑更严格,因此触发了报错。
具体修复步骤:
- 移除
WriteToPubSub之后的"pubsub out print" >> beam.Map(print)变换,因为PubSub写入操作没有输出元素,无法被后续变换处理。 - 如果需要验证输出内容,建议在
WriteToPubSub之前添加打印步骤,示例代码调整如下:groupedRows = denormalised_co_rows | beam.GroupBy(get_hash) | 'ExtractClientIds' >> beam.Map(lambda element: (element[0], [obj['client_id'] for obj in element[1]], element[1])) | "Convert to string" >> beam.Map(encode_as_task) | "Print before PubSub" >> beam.Map(print) | "Write to Pub/Sub" >> beam.io.WriteToPubSub(topic=publish_topic) - 关于权限和窗口选项的疑问澄清:
- 权限问题通常会抛出明确的权限报错(比如PubSub发布权限、CloudSQL连接权限),当前报错不符合权限问题特征,可暂时排除。
- 你的管道属于批处理模式(基于
dummy_input触发数据库读取),无需额外指定窗口选项,窗口主要用于流处理场景。
内容的提问来源于stack exchange,提问作者AnandShiva
相关产品推荐
相关产品推荐

