Airflow如何配置实现上游失败时任务标记为upstream_failed
Airflow DAG 分支触发规则实现方案
要实现你描述的状态判定规则,优先使用Airflow内置能力即可,不需要编写复杂的自定义逻辑。
最优方案(Airflow 2.0+ 适用)
直接将regular_task的trigger_rule参数设置为none_failed_min_one_success即可,该内置规则的判定逻辑完全匹配你的三个要求:
- 任意上游sensor处于
failed或upstream_failed状态时,regular_task会被直接标记为upstream_failed,无需等待其余上游任务执行完成 - 所有上游sensor执行结束、无失败状态,但没有任何sensor执行成功(即两个sensor均为
skipped状态)时,regular_task会被标记为skipped - 所有上游sensor无失败状态,且至少1个sensor执行成功时,
regular_task会正常启动运行
该触发规则是Airflow 2.0版本正式发布的稳定能力,不需要新增任何额外任务,改完参数即可生效,覆盖绝大多数生产环境场景。
低版本兼容方案(Airflow 1.x 适用)
如果你使用的Airflow版本低于2.0,没有内置上述触发规则,可以通过新增轻量状态检查任务的方式实现相同效果,无版本兼容问题:
- 在两个sensor和
regular_task之间新增一个PythonOperator作为状态检查节点,配置该节点的trigger_rule='all_done' - 在检查节点的执行逻辑中拉取两个上游sensor的运行状态,根据状态对应抛出异常或正常返回
- 将
regular_task的触发规则改回默认的all_success即可
参考实现代码:
from airflow.operators.python import PythonOperator, get_current_context from airflow.exceptions import AirflowSkipException, AirflowFailException from airflow.utils.state import State def check_upstream_sensor_status(): context = get_current_context() task_instance = context["ti"] # 拉取两个上游sensor的运行状态 sensor1_state = task_instance.get_state( task_id="branch_1_sensor", dag_id=task_instance.dag_id, execution_date=task_instance.execution_date ) sensor2_state = task_instance.get_state( task_id="branch_2_sensor", dag_id=task_instance.dag_id, execution_date=task_instance.execution_date ) # 任意sensor失败则抛出失败异常,下游regular_task会自动标记为upstream_failed if State.FAILED in (sensor1_state, sensor2_state) or State.UPSTREAM_FAILED in (sensor1_state, sensor2_state): raise AirflowFailException("上游分支sensor存在执行失败") # 两个sensor均被跳过则抛出跳过异常 if sensor1_state == State.SKIPPED and sensor2_state == State.SKIPPED: raise AirflowSkipException("所有上游分支sensor均被跳过,无满足触发条件的分支") # 其余场景(至少一个sensor执行成功)正常返回,触发下游任务运行 return # 定义状态检查任务 status_check = PythonOperator( task_id="upstream_status_check", python_callable=check_upstream_sensor_status, trigger_rule="all_done" ) # 调整DAG依赖关系 branch_1_sensor >> status_check branch_2_sensor >> status_check status_check >> regular_task
原有配置不生效的原因
- 配置
trigger_rule='one_success'时,规则逻辑只会统计上游成功任务的数量,只要上游存在失败任务且最终成功数不足1,任务会直接被标记为skipped,不会识别上游失败的场景 - 配置
trigger_rule='all_done'时,规则只要求所有上游执行结束就触发任务运行,不会区分上游是成功、失败还是跳过状态,因此会出现全跳过场景下任务仍启动的问题
内容的提问来源于stack exchange,提问作者megadarkfriend
相关产品推荐
相关产品推荐

