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

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的校验逻辑更严格,因此触发了报错。

具体修复步骤:

  1. 移除WriteToPubSub之后的"pubsub out print" >> beam.Map(print)变换,因为PubSub写入操作没有输出元素,无法被后续变换处理。
  2. 如果需要验证输出内容,建议在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)
    
  3. 关于权限和窗口选项的疑问澄清:
    • 权限问题通常会抛出明确的权限报错(比如PubSub发布权限、CloudSQL连接权限),当前报错不符合权限问题特征,可暂时排除。
    • 你的管道属于批处理模式(基于dummy_input触发数据库读取),无需额外指定窗口选项,窗口主要用于流处理场景。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 11:56:21