【发布时间】:2018-08-07 13:49:57
【问题描述】:
我有以下带有 3 个任务的 DAG:
start --> special_task --> end
中间的任务可以成功也可以失败,但是end必须始终被执行(假设这是一个干净关闭资源的任务)。为此,我使用了trigger rule ALL_DONE:
end.trigger_rule = trigger_rule.TriggerRule.ALL_DONE
使用它,如果special_task 失败,end 会正确执行。但是,由于 end 是最后一个任务并且成功,因此 DAG 始终标记为 SUCCESS。
如何配置我的 DAG,以便在其中一项任务失败时,将整个 DAG 标记为 FAILED?
重现示例
import datetime
from airflow import DAG
from airflow.operators.bash_operator import BashOperator
from airflow.utils import trigger_rule
dag = DAG(
dag_id='my_dag',
start_date=datetime.datetime.today(),
schedule_interval=None
)
start = BashOperator(
task_id='start',
bash_command='echo start',
dag=dag
)
special_task = BashOperator(
task_id='special_task',
bash_command='exit 1', # force failure
dag=dag
)
end = BashOperator(
task_id='end',
bash_command='echo end',
dag=dag
)
end.trigger_rule = trigger_rule.TriggerRule.ALL_DONE
start.set_downstream(special_task)
special_task.set_downstream(end)
This post 似乎是相关的,但答案不适合我的需要,因为必须执行下游任务end(因此必须执行trigger_rule)。
【问题讨论】:
-
我不知道在 DAG 级别配置此功能的方法。您可以使用任务流来让其他东西传播失败状态,或者使用
on_failure_callback来获得有关失败任务的通知。 -
@JustinasMarozas 实际上,我已经有一个
on_failure_callback来获得通知,但我希望在 Web UI 中将我的 DAG 标记为failed。 -
如果您创建一个虚拟任务并将其设置为
special_task的下游,我希望传播失败。不过,这更像是一种绷带,而不是一种解决方案。 -
@JustinasMarozas 确实,您的解决方案有效,谢谢!但我认为存在一个开箱即用的解决方案,因为它是一个非常常见的用例。但是,对于面临相同问题的人,我会用您的解决方案回答问题,如果没有找到其他解决方案,我会将其标记为答案。感谢您的帮助。
标签: airflow