Airflow技术问询:如何获取上游失败任务列表?
解决方案
你可以用PythonOperator替代原来的DummyOperator,通过自定义Python代码收集并输出所有失败的上游任务列表。核心思路是利用Airflow上下文获取当前DAG运行实例,遍历上游任务的执行状态,筛选出失败任务并打印到标准输出。
完整代码示例
from airflow import DAG from airflow.operators.bash import BashOperator from airflow.operators.python import PythonOperator from airflow.utils.state import State from datetime import datetime def collect_failed_tasks(**context): # 获取当前DAG运行实例 dag_run = context['dag_run'] # 获取当前任务的所有上游任务ID upstream_task_ids = context['task'].upstream_task_ids # 收集失败的任务 failed_tasks = [] for task_id in upstream_task_ids: # 获取对应任务的实例 task_instance = dag_run.get_task_instance(task_id) if task_instance.state == State.FAILED: failed_tasks.append(task_id) # 打印失败任务列表到标准输出 if failed_tasks: print(f"本次DAG运行中失败的上游任务列表: {failed_tasks}") else: print("所有上游任务均执行成功") with DAG( dag_id='failed_task_collector', start_date=datetime(2024, 1, 1), schedule_interval=None, catchup=False ) as dag: # 修正原代码中重复task_id的问题,每个任务需设置唯一ID a = BashOperator(task_id='task_a', bash_command='exit 1') b = BashOperator(task_id='task_b', bash_command='exit 0') c = BashOperator(task_id='task_c', bash_command='exit 1') # 替换原all_success任务为PythonOperator collect_failed = PythonOperator( task_id='collect_failed_tasks', python_callable=collect_failed_tasks, # 设置trigger_rule为all_done,确保所有上游任务结束后无论状态如何都执行 trigger_rule='all_done', provide_context=True ) # 设置任务依赖 [a, b, c] >> collect_failed
关键细节说明
- trigger_rule设置为
all_done:确保不管上游任务成功还是失败,这个收集任务都会执行,才能完整统计所有失败任务。 - 利用Airflow上下文:通过
context参数获取当前DAG运行实例和上游任务ID,进而查询每个任务的执行状态。 - 任务ID唯一性:原代码中三个BashOperator的task_id重复会导致Airflow报错,因此修改为唯一的
task_a、task_b、task_c。
运行后,你可以在该任务的日志中看到标准输出的失败任务列表,对于大型DAG来说,直接查看这个任务的日志就能快速定位所有失败的上游任务。
内容的提问来源于stack exchange,提问作者Vikrant Goel
相关产品推荐
相关产品推荐

