Apache Beam/Dataflow并行任务单文件错误致全流失败排查
Apache Beam写入Firestore单流失败导致全流终止的问题分析与解决
问题描述
用Apache Beam编写代码从GCS并行读取多个文件,写入Firestore不同集合。其中file_2包含重复数据,写入时触发Firestore错误:400 A non-transactional commit may not contain multiple mutations affecting the same entity。但file_2处理失败后,所有并行文件处理流都跟着失败,预期仅file_2对应流报错,怀疑各流未完全独立并行。
测试代码
import apache_beam as beam from apache_beam.options.pipeline_options import PipelineOptions from apache_beam.io.gcp.datastore.v1new.datastoreio import WriteToDatastore from apache_beam.io.gcp.datastore.v1new.types import Key from apache_beam.io.gcp.datastore.v1new.types import Entity import hashlib import json from src.settings import LOCAL_SCHEMAS import os import logging import re from datetime import datetime, timedelta PROJECT = "test-project" class CSVtoDict(beam.DoFn): """Converts line into dictionary""" def process(self, element, headers): element = re.findall(r"[^,"]+|\"[^\"]+\"", element) yield dict(zip(headers, element)) class EntityWrapper(object): """ Create a Cloud firestore entity from the given string. Namespace and project are taken from the parent key. """ def __init__(self, kind, exclude_from_indexes, columns, parent_key=None): self._kind = kind self._exclude_from_indexes = exclude_from_indexes self._parent_key = parent_key self._columns = columns logging.info("Start entities creation") self._expiry_date = datetime.today() + timedelta(days=3) def make_entity(self, content): """Create entity from given string.""" key_columns = list(set(self._columns) - set(self._exclude_from_indexes)) entity_keys = {k: content[k] for k in key_columns} key = Key( [ self._kind, hashlib.sha1(json.dumps(entity_keys).encode("utf-8")).hexdigest(), ], parent=self._parent_key, namespace="test-namespace", ) excluded_indexes = self._exclude_from_indexes excluded_indexes.extend(["expiry_date"]) entity = Entity(key, exclude_from_indexes=excluded_indexes) entity.set_properties({"expiry_date": self._expiry_date}) for index_key in content.keys(): if index_key == "rank": entity.set_properties({index_key: int(content[index_key])}) else: entity.set_properties({index_key: str(content[index_key])}) return entity def dataflow(): """ Upload data from GCS to firestore """ pipeline_options = { "runner": "DirectRunner", "job_name": "firestore-upload-{}".format( datetime.now().strftime("%Y-%m-%d-%H%M%S") ), } options = PipelineOptions.from_dictionary(pipeline_options) p = beam.Pipeline(options=options) for i, source in enumerate(LOCAL_SCHEMAS.keys()): Kind_schema = LOCAL_SCHEMAS[source] input_filename = Kind_schema["input_filename"] input_file_path = os.path.join("dummy_data", input_filename) columns = Kind_schema["columns"] kind_name = Kind_schema["kind_name"] exclude_from_indexes = Kind_schema["exclude_from_indexes"] print("upload {} to {} in firestore".format(input_file_path, kind_name)) ( p | "Reading input file_{}".format(i) >> beam.io.ReadFromText(input_file_path, skip_header_lines=1) | "Converting from csv to dict_{}".format(i) >> beam.ParDo(CSVtoDict(), columns) | "Create entities_{}".format(i) >> beam.Map( EntityWrapper(kind_name, exclude_from_indexes, columns).make_entity ) | "Write entities into firestore_{}".format(i) >> WriteToDatastore(PROJECT) ) result = p.run() result.wait_until_finish() if __name__ == "__main__": dataflow()
测试文件说明
file_1(无重复)
category_id,product_id,rank,seed_product 98354900,53317,6,59596
file_2(含重复)
product_id,category_name,category_id,number_of_purchases,rank 1091056,Womts & Jackets,C_600001530,24357,1 1091052,clo,C_190118298,24357,1 1095256,Women,C_5000298,24357,1 565702,Electcals,C_50001,22304,1 555702,Home Appes,C_192851,22304,1 5655702,Gri& Fryers,C_8000535,22304,1 655702,Sma Cooking Applians,5900861776,22304,1 5417600,test,60000246,122,1 5417600,test,60000246,122,1 5417600,test,60000246,122,1
file_3(无重复)
product_id,sub_category_name,sub_category_id,number_of_views,rank 1091056,Women's Cos & Jackets,C_601530,2357,1 5652,Grills & Frys,C_80035,22304,1 62787,Women's Drses,C_6001506,20521,1 54600,Duvet Cors,C_7003971,19990,1
file_4(无重复)
product_id,number_of_purchases 1091056,24357 56502,22304
问题原因
- DirectRunner全局错误机制:使用的DirectRunner是本地运行器,错误处理为全局模式——一旦Pipeline中某个变换抛出未捕获异常,整个Pipeline会立即终止,不会隔离失败分支。这和Dataflow Runner等分布式Runner的容错逻辑不同,后者会重试失败任务,分支隔离性更强。
- 批量提交失败传播:WriteToDatastore会将多个实体打包成批量请求提交。file_2的重复数据生成相同Key的实体,触发Firestore错误后,该错误向上传播直接终止整个Pipeline,而非仅终止当前文件处理流。
- 无错误处理逻辑:代码未对实体创建、写入的异常做捕获处理,一旦出错就导致全流程崩溃。
解决思路
1. 预处理去重,从根源避免重复实体
针对重复Key导致的错误,在写入前按实体Key去重:
# 在Create entities步骤后添加去重逻辑 ( p | "Reading input file_{}".format(i) >> beam.io.ReadFromText(input_file_path, skip_header_lines=1) | "Converting from csv to dict_{}".format(i) >> beam.ParDo(CSVtoDict(), columns) | "Create entities_{}".format(i) >> beam.Map( EntityWrapper(kind_name, exclude_from_indexes, columns).make_entity ) | "Deduplicate entities_{}".format(i) >> beam.Distinct(lambda entity: entity.key) # 按实体Key去重 | "Write entities into firestore_{}".format(i) >> WriteToDatastore(PROJECT) )
2. 错误分流,隔离失败任务
使用SideOutputs将错误数据单独输出,不影响主流程执行:
class CreateEntityWithErrorHandling(beam.DoFn): def process(self, content, wrapper): try: yield wrapper.make_entity(content) except Exception as e: # 将错误数据输出到SideOutput yield beam.pvalue.TaggedOutput('entity_errors', (content, str(e))) # 替换原Create entities步骤 entity_wrapper = EntityWrapper(kind_name, exclude_from_indexes, columns) create_entity_step = ( p | "Reading input file_{}".format(i) >> beam.io.ReadFromText(input_file_path, skip_header_lines=1) | "Converting from csv to dict_{}".format(i) >> beam.ParDo(CSVtoDict(), columns) | "Create entities with error handling_{}".format(i) >> beam.ParDo(CreateEntityWithErrorHandling(), entity_wrapper).with_outputs('entity_errors', main='entities') ) # 主流程:写入正常实体 create_entity_step.entities | "Write entities into firestore_{}".format(i) >> WriteToDatastore(PROJECT) # 错误分流:将错误数据写入GCS日志 create_entity_step.entity_errors | "Write errors to GCS_{}".format(i) >> beam.io.WriteToText( os.path.join("error_logs", f"errors_{i}.txt") )
3. 切换到分布式Runner
使用Google Cloud Dataflow Runner,其容错机制会自动重试失败任务,不同文件的处理分支运行在独立Worker上,单个分支失败不会导致全流程终止。修改pipeline_options:
pipeline_options = { "runner": "DataflowRunner", "job_name": "firestore-upload-{}".format(datetime.now().strftime("%Y-%m-%d-%H%M%S")), "project": PROJECT, "region": "us-central1", # 替换为你的GCP区域 "temp_location": "gs://your-bucket/temp", # 替换为你的GCS临时路径 }
4. 调整批量提交大小(临时方案)
减小WriteToDatastore的批量大小,降低单个批量失败的影响,但会牺牲写入性能:
| "Write entities into firestore_{}".format(i) >> WriteToDatastore( PROJECT, batch_size=50 # 减小批量大小,默认值更大 )
内容的提问来源于stack exchange,提问作者khalil
相关产品推荐
相关产品推荐

