代码之家  ›  专栏  ›  技术社区  ›  Christopher Beck

分支后的气流任务Operator未正确失败和成功

  •  3
  • Christopher Beck  · 技术社区  · 8 年前

    在我的DAG中,我有一些任务应该只在星期六运行。因此,我使用branchpython运算符在周六的任务和DummyTask之间进行分支。之后,我加入两个分支,并希望运行其他任务。

    工作流如下所示: enter image description here
    在这里,我将dummy3的触发规则设置为 'one_success' 一切都很好。

    我遇到的问题是当branchpythonooperator的上游发生故障时: enter image description here
    分支运算符和分支的状态正确 '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
    

    编辑

    我刚想出一个难看的解决办法: enter image description here
    dummy4表示我实际上需要运行的任务,dummy5只是一个虚拟任务。
    dummy3仍然有触发规则 “一个成功” .

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

    这个解决方案可以让我的DAG按我所希望的方式运行,但我仍然希望没有一些复杂的解决方案。

    2 回复  |  直到 8 年前
        1
  •  2
  •   Géraud    7 年前

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

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


    编辑 :看起来像这样 '没有失败' 当这个问题被问及回答时,触发规则还没有存在:它在2018年11月被添加。

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

        2
  •  4
  •   Alessandro Cosentino    8 年前

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

    它按预期工作,而且可以说它比您的解决方案更干净,因为您没有任何附加的子虚拟运算符。但是,您丢失了平面结构,现在必须放大子图以查看内部结构的详细信息。


    一个更普遍的观察:在对你的DAG进行实验后,我得出结论,气流需要像JoinOperator这样的东西来代替你的Dummy3 operator。让我解释一下。您描述的行为来自这样一个事实,即DAG的成功仅基于最后一个成功(或跳过)的运算符.

    下面的DAG以Success status结尾,是一个支持上述声明的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只在 立即的 父母是成功的,其他的都被跳过,不用使用 trigger_rule 争论。

    或者,解决你面临的问题的方法就是触发规则 all (success | skipped) ,你可以向Dummy3申请。不幸的是,我认为你还不能在气流上创建自定义触发规则。

    编辑 :在这个答案的第一个版本中,我声明触发器规则 one_success 和 all_success 根据成功的程度开火 全部的 达格人的祖先,而不仅仅是直系父母。这与 documentation 事实上,通过以下实验,它是无效的: https://gist.github.com/cosenal/b607825539aa0d308f10f3095e084fac