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

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

问题原因

  1. DirectRunner全局错误机制:使用的DirectRunner是本地运行器,错误处理为全局模式——一旦Pipeline中某个变换抛出未捕获异常,整个Pipeline会立即终止,不会隔离失败分支。这和Dataflow Runner等分布式Runner的容错逻辑不同,后者会重试失败任务,分支隔离性更强。
  2. 批量提交失败传播:WriteToDatastore会将多个实体打包成批量请求提交。file_2的重复数据生成相同Key的实体,触发Firestore错误后,该错误向上传播直接终止整个Pipeline,而非仅终止当前文件处理流。
  3. 无错误处理逻辑:代码未对实体创建、写入的异常做捕获处理,一旦出错就导致全流程崩溃。

解决思路

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 15:07:59