【问题标题】:DAG marked as "success" if one task fails, because of trigger rule ALL_DONE如果一项任务因触发规则 ALL_DONE 而失败,则 DAG 标记为“成功”
【发布时间】: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


【解决方案1】:

我认为这是一个有趣的问题,并花了一些时间弄清楚如何在没有额外的虚拟任务的情况下实现它。这变成了一项多余的任务,但最终结果如下:

这是完整的 DAG:

import airflow
from airflow import AirflowException
from airflow.models import DAG, TaskInstance, BaseOperator
from airflow.operators.bash_operator import BashOperator
from airflow.operators.dummy_operator import DummyOperator
from airflow.operators.python_operator import PythonOperator
from airflow.utils.db import provide_session
from airflow.utils.state import State
from airflow.utils.trigger_rule import TriggerRule

default_args = {"owner": "airflow", "start_date": airflow.utils.dates.days_ago(3)}

dag = DAG(
    dag_id="finally_task_set_end_state",
    default_args=default_args,
    schedule_interval="0 0 * * *",
    description="Answer for question https://stackoverflow.com/questions/51728441",
)

start = BashOperator(task_id="start", bash_command="echo start", dag=dag)
failing_task = BashOperator(task_id="failing_task", bash_command="exit 1", dag=dag)


@provide_session
def _finally(task, execution_date, dag, session=None, **_):
    upstream_task_instances = (
        session.query(TaskInstance)
        .filter(
            TaskInstance.dag_id == dag.dag_id,
            TaskInstance.execution_date == execution_date,
            TaskInstance.task_id.in_(task.upstream_task_ids),
        )
        .all()
    )
    upstream_states = [ti.state for ti in upstream_task_instances]
    fail_this_task = State.FAILED in upstream_states

    print("Do logic here...")

    if fail_this_task:
        raise AirflowException("Failing task because one or more upstream tasks failed.")


finally_ = PythonOperator(
    task_id="finally",
    python_callable=_finally,
    trigger_rule=TriggerRule.ALL_DONE,
    provide_context=True,
    dag=dag,
)

succesful_task = DummyOperator(task_id="succesful_task", dag=dag)

start >> [failing_task, succesful_task] >> finally_

查看由 PythonOperator 调用的_finally 函数。这里有几个关键点:

  1. 使用@provide_session注释并添加参数session=None,以便您可以使用session查询Airflow DB。
  2. 查询当前任务的所有上游任务实例:
upstream_task_instances = (
    session.query(TaskInstance)
    .filter(
        TaskInstance.dag_id == dag.dag_id,
        TaskInstance.execution_date == execution_date,
        TaskInstance.task_id.in_(task.upstream_task_ids),
    )
    .all()
)
  1. 从返回的任务实例中,获取状态并检查State.FAILED 是否在其中:
upstream_states = [ti.state for ti in upstream_task_instances]
fail_this_task = State.FAILED in upstream_states
  1. 执行您自己的逻辑:
print("Do logic here...")
  1. 最后,如果fail_this_task=True:则任务失败:
if fail_this_task:
    raise AirflowException("Failing task because one or more upstream tasks failed.")

最终结果:

【讨论】:

  • 这行得通,但它错误地将“finally”设置为失败,而实际上却没有。如果你能把它标记为上游失败会更好。
【解决方案2】:

正如@JustinasMarozascomment 中解释的那样,解决方案是创建一个虚拟任务,例如:

dummy = DummyOperator(
    task_id='test',
    dag=dag
)

并将其下游绑定到special_task

failing_task.set_downstream(dummy)

因此,DAG 被标记为失败,dummy 任务被标记为upstream_failed

希望有一个开箱即用的解决方案,但等待这个解决方案可以完成工作。

【讨论】:

    【解决方案3】:

    为了扩展 Bas Harenslak 的答案,一个更简单的 _finally 函数可以检查所有任务(不仅是上游任务)的状态:

    def _finally(**kwargs):
        for task_instance in kwargs['dag_run'].get_task_instances():
            if task_instance.current_state() != State.SUCCESS and \
                    task_instance.task_id != kwargs['task_instance'].task_id:
                raise Exception("Task {} failed. Failing this DAG run".format(task_instance.task_id))
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2019-10-05
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2015-06-15
      • 2014-05-05
      • 1970-01-01
      相关资源
      最近更新 更多