Airflow 1.x如何配置父任务失败运行子任务、上游失败时跳过的规则
Airflow 1.x 区分
failed与upstream_failed状态的触发规则实现方案 Airflow 1.x 确实没有内置可以区分failed和upstream_failed状态的Trigger Rule,两种实现方案如下:
方案1:自定义Trigger Rule(推荐)
可以直接修改Airflow核心源码新增自定义Trigger Rule,一次修改全局可用:
- 找到Airflow安装路径下的
airflow/utils/trigger_rule.py文件 - 新增自定义规则常量,同时加入合法规则集合
- 调整状态判断逻辑匹配需求:所有父任务状态为
success/failed时执行子任务,存在任意父任务为upstream_failed/skipped时跳过子任务
修改示例代码:
from airflow.utils.state import State # 新增自定义规则常量 ALL_SUCCESS_OR_FAILED = "all_success_or_failed" # 将新规则加入全局合法Trigger Rule集合 TRIGGER_RULES = { ALL_SUCCESS, ALL_FAILED, ALL_DONE, ONE_SUCCESS, ONE_FAILED, NONE_FAILED, NONE_SKIPPED, ALL_SKIPPED, DUMMY, ALL_SUCCESS_OR_FAILED # 新增自定义规则 } # 新增对应判断逻辑(可直接整合到调度器的状态判断流程中) def judge_custom_trigger(upstream_states: set) -> bool: skip_states = {State.UPSTREAM_FAILED, State.SKIPPED} allowed_run_states = {State.SUCCESS, State.FAILED} if any(s in skip_states for s in upstream_states): return False return all(s in allowed_run_states for s in upstream_states)
修改完成后重启Airflow的scheduler和webserver服务,即可在DAG的任务定义中直接指定trigger_rule="all_success_or_failed"使用。
方案2:无侵入轻量替代方案
如果没有权限修改Airflow核心源码,可以通过ShortCircuitOperator实现相同逻辑,无需改动全局配置:
- 给子任务前置一个状态判断任务,设置该判断任务的
trigger_rule="all_done",保证无论上游任务状态如何,判断逻辑都能执行 - 在判断逻辑中获取所有上游任务的状态,按需求返回布尔值控制下游任务是否执行
DAG示例代码:
from airflow.models import DAG from airflow.operators.python_operator import ShortCircuitOperator, PythonOperator from airflow.utils.state import State from datetime import datetime def check_upstream_status(**context): ti = context["ti"] dag_run = ti.get_dagrun() # 获取所有上游任务的状态 upstream_states = [ dag_run.get_task_instance(task_id).state for task_id in ti.task.upstream_task_ids ] skip_states = {State.UPSTREAM_FAILED, State.SKIPPED} allowed_states = {State.SUCCESS, State.FAILED} if any(s in skip_states for s in upstream_states): return False return all(s in allowed_states for s in upstream_states) # DAG定义 with DAG( dag_id="custom_trigger_demo", start_date=datetime(2024, 1, 1), schedule_interval=None, catchup=False ) as dag: # 示例父任务 parent_task = PythonOperator( task_id="parent_task", python_callable=lambda: 1 / 0 # 模拟任务运行失败场景 ) # 上游状态判断任务 upstream_check = ShortCircuitOperator( task_id="upstream_check", python_callable=check_upstream_status, provide_context=True, trigger_rule="all_done" ) # 实际要执行的子任务 child_task = PythonOperator( task_id="child_task", python_callable=lambda: print("子任务正常执行") ) parent_task >> upstream_check >> child_task
内容的提问来源于stack exchange,提问作者elaspog
相关产品推荐
相关产品推荐

