Airflow 1.8中all_done触发规则在祖辈节点状态为None时失效的解决方法
我在Airflow 1.8版本里踩过一模一样的坑,太懂这种上游祖辈失败导致清理任务挂死的痛苦了!咱们先把问题根源理清楚,再给你几个可行的解决方案:
问题根源
Airflow 1.8的all_done触发规则有个关键局限:它只检查直接父节点的状态,不会递归遍历所有祖辈节点。当你的祖辈任务失败时,中间层级的父任务(比如load_att_wv2_m1bs)会因为上游未完成而被标记为State.None(而非State.Upstream_Failed),这时候all_done会判定父任务还没"完成",导致下游的tmp_cleanup永远不会进入排队状态,一直卡在None。
而直接父节点失败时能正常工作,是因为直接父节点的状态会变成State.failed(属于Airflow定义的finished_states),all_done会认为父任务已经完成,所以触发下游的清理任务。
解决方案
方案1:自定义递归版的'all_done'触发规则
这是最彻底的解决办法,我们可以写一个自定义触发规则,递归检查所有上游(包括祖辈)任务是否都处于完成状态(不管成功、失败还是跳过):
from airflow.utils.trigger_rule import BaseTriggerRule from airflow.utils.state import State class AllDoneRecursive(BaseTriggerRule): name = 'all_done_recursive' def is_triggered(self, task_instance, dag_run, task): # 递归遍历所有上游任务,确认都已完成 def check_all_upstream(t): # 没有上游任务,直接返回True if not t.upstream_list: return True for upstream_task in t.upstream_list: ti = dag_run.get_task_instance(upstream_task) # 检查当前上游任务是否处于完成状态 if ti.state not in State.finished_states: return False # 递归检查该上游任务的所有上游 if not check_all_upstream(upstream_task): return False return True return check_all_upstream(task) # 注册自定义触发规则 BaseTriggerRule.register_trigger_rule(AllDoneRecursive())
把这段代码放到你的DAG定义文件里,然后给tmp_cleanup任务设置触发规则:
tmp_cleanup = BashOperator( task_id='tmp_cleanup', bash_command='rm -rf /tmp/your_temp_dir', trigger_rule='all_done_recursive', dag=dag )
方案2:临时Workaround(无需修改代码结构)
如果暂时不想写自定义规则,可以给中间层级的父任务也设置all_done触发规则。比如给load_att_wv2_m1bs加上trigger_rule='all_done',这样当祖辈任务失败时,load_att_wv2_m1bs会被触发执行(哪怕执行失败),它的状态会变成State.failed(属于完成状态),此时tmp_cleanup的直接父节点已经"完成",原生的all_done规则就会正常触发清理任务。
方案3:升级Airflow版本(长期建议)
Airflow在后续版本(比如1.10+)里修复了这个触发规则的逻辑,all_done会正确处理祖辈节点失败的场景。如果你的环境允许升级,这是一劳永逸的办法。
内容的提问来源于stack exchange,提问作者7yl4r

