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

解决Apache Beam Dataflow管道中的'PBegin object has no attribute windowing'错误

问题:Apache Beam Dataflow作业中出现AttributeError: 'PBegin' object has no attribute 'windowing'错误

我正在开发一个Apache Beam Dataflow作业,用于将经过富化处理的流数据存储至Firestore。若处理过程中出现事件数据遗漏,我会将这些记录推送至BigQuery表中。但运行时始终遇到错误:AttributeError: 'PBegin' object has no attribute 'windowing'。

相关代码片段

class EnrichEventsWithFirestore(beam.DoFn):
    def __init__(self, database):
        self.database = database
        self.client = None

    def setup(self):
        if not self.client:
            self.client = firestore.Client(project='temp_myself', database=self.database)

    def process(self, element):
        list_ids = []
        for ele in element:
            Id = ele[0]['Id']
            list_ids.append(self.client.document('cust_add', str(Id)))
        docs = self.client.get_all(references=list_ids)
        Id_map = {}
        for doc in docs:
            doc_dict = doc.to_dict()
            if doc_dict:
                Id_map[doc_dict['Id']] = {
                    "c_num": doc_dict.get('c_num', "NA"),
                    "ct_num": doc_dict.get("ct_num", "NA")
                }
        enriched_data = []
        for ele in element:
            tran_data, item_data = ele
            Id = int(tran_data['Id'])
            filename = str(tran_data['filename'])
            tdatetime = tran_data['tdatetime']
           
            if Id in Id_map:
                ct_num = Id_map[Id]['ct_num']
                c_num = Id_map[Id]['c_num']
            else:
                yield beam.pvalue.TaggedOutput('discarded',
                    [{'Id': Id, 'filename': filename, 'tdatetime': tdatetime}]
                )
                continue

            tran_data["ct_num"] = ct_num
            tran_data["c_num"] = c_num
            enriched_data.append((tran_data, item_data))

        yield enriched_data

# Pipeline Code
def run(pipeline_args):
    pipeline_options = PipelineOptions(pipeline_args, save_main_session=True)
    pipeline_options.view_as(StandardOptions).streaming = True
    pipeline = beam.Pipeline(options=pipeline_options)

    if pipeline:
        results = (
            pipeline
            | "Process File Firestore" >> beam.FlatMap(parseForFirestore, campaign__data)
            | "Window and Batch Data Firestore" >> beam.WindowInto(
                beam.window.FixedWindows(5),
                trigger=Repeatedly(AfterAny(AfterCount(100), AfterProcessingTime(5))),
                accumulation_mode=AccumulationMode.DISCARDING,
                allowed_lateness=600
            )
            | "Group Data Firestore" >> beam.GroupBy(lambda s: assign_random_int())
            | "Extract Values Firestore" >> beam.Values()
            | "Enrich Events with Firestore Data" >> beam.ParDo(
                EnrichEventsWithFirestore('temp_myself')
            ).with_outputs('discarded', main='not_discarded')
        )

        results.not_discarded | "Push Data to Firestore" >> beam.ParDo(InsertToFireStoreInBatches('temp_myself'))
        results.discarded | "Write to BQ" >> beam.io.WriteToBigQuery(
            'bq_table',
            create_disposition="CREATE_IF_NEEDED",
            write_disposition="WRITE_APPEND",
            ignore_unknown_columns=True,
            method="STREAMING_INSERTS",
            schema=tab_schema
        )
    
    pipeline.run()

错误原因及修复方案

这个错误的核心原因是TaggedOutput输出的数据没有携带窗口信息,而后续的WriteToBigQuery(使用STREAMING_INSERTS模式)需要处理带窗口的数据流,当数据没有窗口元数据时就会抛出该错误。具体问题出在EnrichEventsWithFirestore的process方法中:

  • 问题点1:TaggedOutput输出的是列表格式,而非单个元素,导致窗口上下文丢失;主输出将批量数据作为单个列表yield,同样破坏了元素的窗口元数据。
  • 问题点2:BigQuery的STREAMING_INSERTS模式依赖数据流的窗口信息来处理流数据,无窗口元数据的输入会触发该AttributeError。

修复步骤:

  1. 调整TaggedOutput输出格式:直接输出单个字典对象,而非列表,确保保留原始元素的窗口上下文
  2. 调整主输出方式:逐个yield富化后的元素,而非将整个列表作为单个元素输出

修改后的EnrichEventsWithFirestore类代码:

class EnrichEventsWithFirestore(beam.DoFn):
    def __init__(self, database):
        self.database = database
        self.client = None

    def setup(self):
        if not self.client:
            self.client = firestore.Client(project='temp_myself', database=self.database)

    def process(self, element):
        list_ids = []
        for ele in element:
            Id = ele[0]['Id']
            list_ids.append(self.client.document('cust_add', str(Id)))
        docs = self.client.get_all(references=list_ids)
        Id_map = {}
        for doc in docs:
            doc_dict = doc.to_dict()
            if doc_dict:
                Id_map[doc_dict['Id']] = {
                    "c_num": doc_dict.get('c_num', "NA"),
                    "ct_num": doc_dict.get("ct_num", "NA")
                }
        
        for ele in element:
            tran_data, item_data = ele
            Id = int(tran_data['Id'])
            filename = str(tran_data['filename'])
            tdatetime = tran_data['tdatetime']
           
            if Id in Id_map:
                ct_num = Id_map[Id]['ct_num']
                c_num = Id_map[Id]['c_num']
                tran_data["ct_num"] = ct_num
                tran_data["c_num"] = c_num
                # 逐个输出富化元素,保留窗口上下文
                yield (tran_data, item_data)
            else:
                # 输出单个字典,而非列表,保留窗口上下文
                yield beam.pvalue.TaggedOutput('discarded',
                    {'Id': Id, 'filename': filename, 'tdatetime': tdatetime}
                )

额外检查点:

  • 确保GroupBy和Values操作未破坏数据的窗口信息,当前窗口配置是合规的
  • 若需批量处理主数据,可后续通过beam.BatchElements操作实现,避免在DoFn中直接输出批量列表

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 21:58:10