Airflow分支任务异常:返回的end_task被标记为跳过的原因
问题
使用BranchPythonOperator做分支任务判断,日志明确显示返回了end_task,但该任务并未被标记为成功,反而被跳过。首次使用分支操作符,参考示例编写的代码逻辑看似没问题,但返回的目标任务还是出现被跳过的情况。
代码实现
END = EmptyOperator(task_id="end_task") def missing_check(): SQL = f""" ... """ big_query_hook = BigQueryHook(gcp_conn_id=BIGQUERY_STG_CONN, use_legacy_sql=False) df = big_query_hook.get_pandas_df(SQL, dialect="standard") if any(df['num_missing_ids'] != 0): return 'diff_email_notif' else: return 'end_task' missing_data_check = BranchPythonOperator( task_id='missing_data_check', python_callable=missing_check, ) START >> processing >> fact_cx_daily_log >> missing_data_check missing_data_check >> email_notif >> END missing_data_check >> END
相关日志
[2023-10-11, 12:34:26 UTC] {python.py:183} INFO - Done. Returned value was: end_task [2023-10-11, 12:34:26 UTC] {python.py:216} INFO - Branch callable return end_task [2023-10-11, 12:34:26 UTC] {skipmixin.py:164} INFO - Following branch end_task [2023-10-11, 12:34:26 UTC] {skipmixin.py:224} INFO - Skipping tasks ['diff_email_notif'] [2023-10-11, 12:34:26 UTC] {taskinstance.py:1345} INFO - Marking task as SUCCESS. dag_id=re_cx_daily_log_v3, task_id=missing_data_check, execution_date=20231011T123255, start_date=20231011T123423, end_date=20231011T123426 [2023-10-11, 12:34:27 UTC] {local_task_job_runner.py:225} INFO - Task exited with return code 0 [2023-10-11, 12:34:27 UTC] {taskinstance.py:2653} INFO - 0 downstream tasks scheduled from follow-on schedule check
解决方案
问题本质:
END任务存在两条上游依赖路径:missing_data_check >> email_notif >> END和missing_data_check >> END。当分支返回end_task时,Airflow会跳过email_notif任务,但默认情况下END需要等待所有上游任务完成(包括被跳过的email_notif),此时END会被标记为跳过状态,因为它的上游有任务未成功执行。修复步骤:
- 给
END任务添加trigger_rule="none_failed_min_one_success"配置,该规则表示只要至少有一条上游路径成功完成,END就会执行并标记为成功,同时忽略被跳过或未执行的上游任务。修改后的END定义如下:END = EmptyOperator(task_id="end_task", trigger_rule="none_failed_min_one_success") - 若不需要
email_notif指向END,也可以调整DAG结构,移除email_notif >> END这条依赖,但这种方式可能不符合你的流程设计,因此更推荐第一种方案。
- 给
内容的提问来源于stack exchange,提问作者random student
相关产品推荐
相关产品推荐

