如何利用BranchPythonOperator根据Airflow前置Task_1状态分支任务?
解决BranchPythonOperator根据前置任务状态分支的问题
你之前的代码核心问题在于:dagrun.get_task_instances('Task_1')返回的是TaskInstance对象的列表,直接访问.state会抛出属性错误,因为列表没有该属性。以下是修正后的可行方案:
修正后的状态判断函数
直接通过dagrun.get_task_instance()获取单个TaskInstance对象,而非列表:
def check_task_status(**context): dagrun = context["dag_run"] # 获取指定任务的单个实例对象 task_instance = dagrun.get_task_instance(task_id='Task_1') # 根据状态返回对应任务ID return 'Task_2' if task_instance.state == 'success' else 'Task_3'
完整DAG示例
确保分支任务的上游关联到Task_1,且trigger_rule='all_done'配置正确(保证无论Task_1成功/失败,分支任务都会执行):
from airflow import DAG from airflow.operators.python import BranchPythonOperator from airflow.operators.dummy import DummyOperator from datetime import datetime def check_task_status(**context): dagrun = context["dag_run"] task_instance = dagrun.get_task_instance(task_id='Task_1') return 'Task_2' if task_instance.state == 'success' else 'Task_3' with DAG( dag_id='task_status_branch_dag', start_date=datetime(2024, 1, 1), schedule_interval=None, catchup=False ) as dag: Task_1 = DummyOperator(task_id='Task_1') branch_task_1 = BranchPythonOperator( task_id='Check_Task_Status', python_callable=check_task_status, trigger_rule='all_done', provide_context=True, dag=dag ) Task_2 = DummyOperator(task_id='Task_2') Task_3 = DummyOperator(task_id='Task_3') # 设置任务依赖 Task_1 >> branch_task_1 >> [Task_2, Task_3]
关键注意事项
- 必须将
Check_Task_Status的上游设置为Task_1,确保DAG运行时能正确关联到对应的任务实例 trigger_rule='all_done'是必要配置,否则Task_1失败时分支任务不会触发
内容的提问来源于stack exchange,提问作者AlmostThere
相关产品推荐
相关产品推荐

