Airflow分支DAG异常:跳过t1时t2、t3也被跳过的问题排查
解决Airflow分支DAG中跳过t1后t2、t3不执行的问题
问题原因
你的代码里t2的上游同时依赖branch_task和t1,Airflow默认触发规则是all_success——必须所有上游任务执行成功,t2才会启动。当分支逻辑选择跳过t1直接走t2时,t1处于未执行状态,不满足all_success的条件,导致t2和t3都无法触发执行。
修复方案
给t2任务设置trigger_rule='one_success',这个规则表示只要任意一个上游任务执行成功,t2就会启动。同时确保分支逻辑的返回值正确。
修改后的完整代码:
from airflow import DAG from airflow.operators.python import PythonOperator, BranchPythonOperator from datetime import datetime def f0(): print("this is function0") return False # return True def f1(): print("this is function1") def f2(): print("this is function2") def f3(): print("this is function3") def branch_chooser(**kwargs): # 根据t0的返回值选择分支:True走t2,False走t1 return 't2' if kwargs['ti'].xcom_pull(task_ids='t0') else 't1' with DAG( 'TEST', default_args={ 'depends_on_past': False, 'retries': 1, }, description='test process', schedule_interval=None, start_date=datetime(2023, 1, 1), catchup=False, ) as dag: t0 = PythonOperator( task_id='t0', python_callable=f0, provide_context=True, dag=dag, ) branch_task = BranchPythonOperator( task_id='Branch_condition', python_callable=branch_chooser, provide_context=True, dag=dag, ) t1 = PythonOperator( task_id='t1', python_callable=f1, dag=dag, ) t2 = PythonOperator( task_id='t2', python_callable=f2, trigger_rule='one_success', # 关键修改:设置触发规则为任意上游成功即可执行 dag=dag, ) t3 = PythonOperator( task_id='t3', python_callable=f3, dag=dag, ) # 保持原有的依赖关系 t0 >> branch_task branch_task >> [t1, t2] t1 >> t2 t2 >> t3
逻辑验证
- 当
t0返回True时:branch_task选择t2,t2的上游branch_task成功,触发t2执行,之后执行t3,t1被跳过 - 当
t0返回False时:branch_task选择t1,t1执行成功后触发t2,再执行t3
内容的提问来源于stack exchange,提问作者Dr.Dan
相关产品推荐
相关产品推荐

