【问题标题】:Airflow task after BranchPythonOperator does not fail and succeed correctlyBranchPythonOperator 之后的 Airflow 任务不会失败并正确成功
【发布时间】:2018-08-03 03:07:39
【问题描述】:

在我的 DAG 中,我有一些只应在星期六运行的任务。因此,我使用 BranchPythonOperator 在星期六的任务和 DummyTask 之间进行分支。之后,我加入了两个分支并想运行其他任务。

工作流程如下所示:
这里我将dummy3的触发规则设置为'one_success',一切正常。

我遇到的问题是 BranchPythonOperator 上游出现故障时:
BranchPythonOperator 和分支正确的状态为'upstream_failed',但加入分支的任务变为'skipped',因此整个工作流显示'success'

我尝试使用'all_success' 作为触发规则,然后如果出现故障则整个工作流程都会失败,但如果没有任何故障,则跳过 dummy3。

我也试过'all_done'作为触发规则,如果没有失败,它会正常工作,但如果有失败dummy3仍然会被执行。

我的测试代码如下所示:

from datetime import datetime, date
from airflow import DAG
from airflow.operators.python_operator import BranchPythonOperator, PythonOperator
from airflow.operators.dummy_operator import DummyOperator

dag = DAG('test_branches',
          description='Test branches',
          catchup=False,
          schedule_interval='0 0 * * *',
          start_date=datetime(2018, 8, 1))


def python1():
    raise Exception('Test failure')
    # print 'Test success'


dummy1 = PythonOperator(
    task_id='python1',
    python_callable=python1,
    dag=dag
)


dummy2 = DummyOperator(
    task_id='dummy2',
    dag=dag
)


dummy3 = DummyOperator(
    task_id='dummy3',
    dag=dag,
    trigger_rule='one_success'
)


def is_saturday():
    if date.today().weekday() == 6:
        return 'dummy2'
    else:
        return 'today_is_not_saturday'


branch_on_saturday = BranchPythonOperator(
    task_id='branch_on_saturday',
    python_callable=is_saturday,
    dag=dag)


not_saturday = DummyOperator(
    task_id='today_is_not_saturday',
    dag=dag
)

dummy1 >> branch_on_saturday >> dummy2 >> dummy3
branch_on_saturday >> not_saturday >> dummy3

编辑

我刚刚想出了一个丑陋的解决方法:
dummy4 代表我实际需要运行的任务,dummy5 只是一个假人。
dummy3 仍然有触发规则'one_success'

现在 dummy3 和 dummy4 在上游没有失败的情况下运行,如果当天不是星期六 dummy5 '运行',如果当天是星期六则跳过,这意味着 DAG 在两种情况下都被标记为成功。
如果上游出现故障,则跳过 dummy3 和 dummy4,将 dummy5 标记为'upstream_failed',并将 DAG 标记为失败。

这种解决方法可以让我的 DAG 按我的意愿运行,但我仍然更喜欢没有一些 hacky 解决方法的解决方案。

【问题讨论】:

    标签: python airflow


    【解决方案1】:

    将 dummy3 的触发规则设置为 'none_failed' 将使其在任何情况下都以预期状态结束。

    https://airflow.apache.org/concepts.html#trigger-rules


    编辑:在询问和回答此问题时,似乎此 'none_failed' 触发规则尚不存在:它于 2018 年 11 月添加

    https://github.com/apache/airflow/pull/4182

    【讨论】:

    • 这是对气流的一个很好的新增功能,将使分支更清洁。
    【解决方案2】:

    您可以使用的一种解决方法是将 DAG 的第二部分放在 SubDAG 中,就像我在以下代码中说明您的示例一样:https://gist.github.com/cosenal/cbd38b13450b652291e655138baa1aba

    它按预期工作,并且可以说比您的解决方法更清洁,因为您没有任何额外的辅助虚拟运算符。但是,您失去了平面结构,现在您必须放大 SubDag 才能查看内部结构的详细信息。


    更一般的观察:在对您的 DAG 进行试验后,我得出的结论是,Airflow 需要类似 JoinOperator 的东西来替换您的 Dummy3 运算符。让我解释。您描述的行为来自这样一个事实,即 DAG 的成功仅基于最后一个运算符成功(或跳过!)。

    以下以“Success”状态结尾的 DAG 是支持上述声明的 MWE。

    def python1():
        raise Exception('Test failure')
    
    dummy1 = PythonOperator(
        task_id='python1',
        python_callable=python1,
        dag=dag
    )
    
    dummy2 = DummyOperator(
        task_id='dummy2',
        dag=dag,
        trigger_rule='one_success'
    )
    
    dummy1 >> dummy2
    

    如果有一个 JoinOperator 只在 immediate 父母之一成功并且所有其他父母都被跳过时触发,那就太酷了,而不必使用trigger_rule 参数。

    或者,可以解决您面临的问题的方法是触发规则all (success | skipped),您可以将其应用于 Dummy3。不幸的是,我认为您还不能在 Airflow 上创建自定义触发规则。

    编辑:在这个答案的第一个版本中,我声称触发规则 one_successall_success 根据所有的祖先的成功程度触发DAG 中的运算符,而不仅仅是直系父母。这与documentation 不匹配,实际上它被以下实验无效:https://gist.github.com/cosenal/b607825539aa0d308f10f3095e084fac

    【讨论】:

    • 非常感谢将 DAG 的第二部分放在 SubDag 中的建议。我可能会考虑使用它,尽管我在实际工作流程中使用的大多数 Operator 都是 SubDag,因此 SubDag 中的 SubDag 可能不是最好的结构。
    • 另外,你确定你对触发规则的定义是正确的吗?在文档中它说:trigger_rule 的默认值为 all_success,可以定义为“当所有直接上游任务都成功时触发此任务”。此外,我完全同意应该有更多的触发规则、JoinOperator 或一种方法来定义如果任务失败或最后一个任务被跳过,您希望整个工作流失败。
    • @ChristopherBeck 你是对的,在他们谈论“直接上游”的文档中,事实上我做了一个实验,使我的主张无效:gist.github.com/cosenal/b607825539aa0d308f10f3095e084fac我将相应地编辑我的答案。谢谢!
    猜你喜欢
    • 2019-07-21
    • 2019-01-13
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-06-08
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多