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

多步骤工作流失败处理器编排:任务失败时返回数据的最佳实践

Airflow「标记任务失败+传递数据」场景最佳实践

问题背景

现有测试数据与Airflow工作流,目标是让processing_group任务组输出各步骤失败项及全流程成功项的列表,且每个失败处理器需接收item及步骤数据执行回滚操作。但Airflow任务无法同时抛出异常与返回值,需实现「标记任务失败+传递数据」的逻辑。

测试数据

item_input = [
    {"value": 1, "fail_step": 1},
    {"value": 2, "fail_step": 1},
    {"value": 3, "fail_step": 2},
    {"value": 4, "fail_step": 2},
    {"value": 5, "fail_step": 3},
    {"value": 6, "fail_step": 3},
    {"value": 7, "fail_step": None},
    {"value": 8, "fail_step": None},
    {"value": 9, "fail_step": None},   
]

现有工作流代码

with DAG(dag_id="example_dag"):

    @task_group
    def processing_group(item):
        @task
        def step_1(item):
            if item["fail_step"] == 1:
                raise ValueError("Failed at step 1")
            return item
    
        @task(trigger_rule="all_failed")
        def step_1_fail_handler(step_1_res, item):
            print("Step 1 fail:", item)
            return item

        # rest steps
    
        step1_res = step_1(item)
        step1_failed = step_1_fail_handler(step1_res, item)
    
        # rest of wiring
    
        return {
            "step1_failed": step1_failed,
            "step2_failed": step2_failed,
            "step3_failed": step3_failed,
            "step3_success": step3_res,
        }

    processed = processing_group.expand(item=item_input)

    @task
    def simple_buffer(results):
        return results
    
    buffered = simple_buffer(processed)

最佳实践方案

1. 状态返回+分支判断(推荐)

放弃直接在步骤任务中抛异常,改为返回包含执行状态和数据的结构化结果,通过分支任务判断状态流向,既保留数据传递能力,又能在UI中通过分支任务的执行状态区分成功/失败场景,同时失败处理器可直接拿到完整数据做回滚。

修改后的核心代码示例:

@task_group
def processing_group(item):
    @task
    def step_1(item):
        if item["fail_step"] == 1:
            # 返回失败状态+数据+错误信息,不抛异常
            return {"status": "failed", "item": item, "error": "Failed at step 1"}
        return {"status": "success", "item": item}

    @task.branch
    def check_step1(result):
        # 根据状态分支到失败处理器或下一步骤
        if result["status"] == "failed":
            return "processing_group.step_1_fail_handler"
        return "processing_group.step_2"

    @task
    def step_1_fail_handler(result):
        # 直接获取失败item,执行回滚操作
        failed_item = result["item"]
        print(f"Step 1 回滚处理: {failed_item}")
        # 返回失败项信息用于最终汇总
        return {"step": 1, "item": failed_item}

    @task
    def step_2(result):
        item = result["item"]
        if item["fail_step"] == 2:
            return {"status": "failed", "item": item, "error": "Failed at step 2"}
        return {"status": "success", "item": item}

    # 分支流向控制
    step1_result = step_1(item)
    step1_branch = check_step1(step1_result)
    step1_fail_out = step_1_fail_handler(step1_result)
    step2_result = step_2(step1_result)

    # 聚合当前item的最终结果
    @task
    def aggregate(step1_fail=None, step2_fail=None, step3_fail=None, success=None):
        output = {}
        if step1_fail:
            output["step1_failed"] = step1_fail
        elif step2_fail:
            output["step2_failed"] = step2_fail
        elif step3_fail:
            output["step3_failed"] = step3_fail
        elif success:
            output["success"] = success["item"]
        return output

    # 根据分支结果传递数据到聚合任务
    final_aggregate = aggregate(
        step1_fail=step1_fail_out,
        # 其他步骤的结果传递逻辑类似
    )

    return final_aggregate

这种方式逻辑清晰,避免XCom依赖,完美支持expand动态任务场景,每个item的处理链路独立可控。

2. XCom预推送+异常抛出(需标记任务失败时使用)

如果必须在Airflow UI中标记步骤任务为失败状态(红色标识),可以先将数据推送到XCom,再抛出异常。失败处理器通过XCom获取数据执行回滚。

核心代码示例:

from airflow.models import TaskInstance
from airflow.utils.session import provide_session

@task_group
def processing_group(item):
    @task
    def step_1(item, ti: TaskInstance, session=None):
        if item["fail_step"] == 1:
            # 先将失败item推送到XCom
            ti.xcom_push(key="failed_item", value=item, session=session)
            raise ValueError("Failed at step 1")
        return item

    @task(trigger_rule="all_failed")
    def step_1_fail_handler(ti: TaskInstance):
        # 从当前任务实例的XCom中拉取失败数据
        failed_item = ti.xcom_pull(key="failed_item", task_ids="processing_group.step_1")
        print(f"Step 1 回滚处理: {failed_item}")
        return {"step": 1, "item": failed_item}

    step1_res = step_1(item)
    step1_fail_out = step_1_fail_handler(step1_res)

    # 后续步骤及聚合逻辑类似方案1
    return step1_fail_out

注意:expand动态任务场景下,每个任务实例有独立的XCom上下文,不会出现数据冲突,但需注意XCom的key命名唯一性,避免不同步骤的XCom互相覆盖。

3. 最终结果汇总

在simple_buffer任务中,可以统一整理所有item的处理结果,生成各步骤失败列表和成功列表:

@task
def simple_buffer(results):
    step1_failed = [r["step1_failed"]["item"] for r in results if "step1_failed" in r]
    step2_failed = [r["step2_failed"]["item"] for r in results if "step2_failed" in r]
    step3_failed = [r["step3_failed"]["item"] for r in results if "step3_failed" in r]
    success_items = [r["success"] for r in results if "success" in r]
    
    return {
        "step1_failed": step1_failed,
        "step2_failed": step2_failed,
        "step3_failed": step3_failed,
        "success_items": success_items
    }

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.11 22:13:20