多步骤工作流失败处理器编排:任务失败时返回数据的最佳实践
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
相关产品推荐
相关产品推荐

