Airflow中按条件执行的任务
我在尝试编写一个有向无环图(DAG),让它在条件成立时执行另一项任务。简化版本如下:
to_be_triggered = EmptyOperator(task_id="to_be_triggered")
@task.branch()
def trigger_dag(**kwargs):
config = kwargs.get("dag_run_config")
if config.get("run_trigger") is True:
return ["to_be_triggered"]
return None
with DAG("example") as dag:
dag_run_config = {
"run_trigger": True
}
t0 = trigger_dag(dag_run_config=dag_run_config)
t1 = EmptyOperator(task_id="end", trigger_rule=TriggerRule.ONE_SUCCESS)
t0 >> t1
因此如果配置中的 run_trigger 变量为真,我想有条件地执行 to_be_triggered。我之所以做不到,是因为 branch_task_ids 只能包含有效的task_id,而且出于某种原因,to_be_triggered 是无效的:
Following branch {'to_be_triggered'}
Task failed with exception
AirflowException: 'branch_task_ids' must contain only valid task_ids. Invalid tasks found: {'to_be_triggered'}
据我从谷歌上了解,这通常是因为某个任务处于一个任务组中,需要用组ID来指定,但我这里没有任务组。有人知道任务组是否在某处被隐式设置,或者 to_be_triggered 无效还有其他可能的原因吗?
解决方案
关键修复点是:
-
在DAG的上下文中定义
to_be_triggered。 -
让它直接位于分支任务的下游:
t0 >> [to_be_triggered, end]
- 如果你希望DAG能顺利地继续到一个结束任务,请为假分支返回
"end",而不是None。
你也可以返回 None;Airflow支持用它来跳过所有下游任务。但就你的结构而言,返回 "end" 更清晰,因为你明确地把“啥也不做,顺利结束”这个分支建模出来。
站内所有文章版权归属LeftHeroAI导航站,无授权禁止任何主体转载、抄袭、复制内容,亦不得私自架设镜像站点。一经侵权,本站将通过法律途径追责。