解决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。
修复步骤:
- 调整TaggedOutput输出格式:直接输出单个字典对象,而非列表,确保保留原始元素的窗口上下文
- 调整主输出方式:逐个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
相关产品推荐
相关产品推荐

