【问题标题】:How to create a multiple conditional tasks in Airflow?如何在 Airflow 中创建多个条件任务?
【发布时间】:2022-08-19 16:57:24
【问题描述】:

要求是当任务 A 失败时,我需要触发任务 D1。同样,当任务 B 失败时,需要触发任务 D2。

我已将其编码如下。使用此代码,当 Task-A 失败时,Task D1 和 Task D2 都会被触发。如何解决这个问题?

def do_op1_work(**kwargs):
    x = kwargs[\'dag_run\'].conf.get(\'x\')
    log.info(\'x: \' + str(x))
    if x == 0:
        raise ValueError(\'Manual Exception\')


def do_op2_work(**kwargs):
    y = kwargs[\'dag_run\'].conf.get(\'y\')
    log.info(\'y: \' + str(y))
    if y == 0:
        raise ValueError(\'Manual Exception\')


with DAG(dag_id=\'fulfill_uv\', schedule_interval=None, default_args=default_args, catchup=False) as dag:
    op1 = PythonOperator(task_id=\'A\', python_callable=do_op1_work, provide_context=True)
    op2 = PythonOperator(task_id=\'B\', python_callable=do_op2_work, provide_context=True)
    op3 = DummyOperator(task_id=\'C\')
    op4 = DummyOperator(task_id=\'D1\', trigger_rule=\'all_failed\')
    op5 = DummyOperator(task_id=\'D2\', trigger_rule=\'all_failed\')
    op6 = DummyOperator(task_id=\'E\', trigger_rule=\'one_success\')

    op2.set_upstream(op1)
    op3.set_upstream(op2)
    op4.set_upstream(op1)
    op5.set_upstream(op2)
    op6.set_upstream(op3)
    op6.set_upstream(op4)
    op6.set_upstream(op5)
  • 如果任务 A 成功,它运行任务 B,但如果它失败,它运行任务 D?

标签: airflow


【解决方案1】:

当任务 A 失败时,任务 B 的状态为upstream_failed,被认为是失败,此时 D2 将被执行。

为了简化你的dag的逻辑,绕过这个问题,你可以创建两个BranchPythonOperator

  1. 获取任务 A 的状态,如果是 failed,则运行 D1,如果是 succeeded,则运行 B
  2. 第二个获取任务 B 的状态,如果是 failed,则运行 D2,如果是 succeeded,则运行 C

    获取状态:

    def get_state(task_id, **context):
        return context["dag_run"].get_task_instance(task_id).state
    

    对于分支运营商:

    from airflow.operators.python_operator import BranchPythonOperator
    
    def branch_func(task_to_check, task_to_run_on_success, task_to_run_on_fail, **context):
        task_state = get_state(task_to_check, **context)
        if task_state == "succeeded":
            return task_to_run_on_success
        else:
            return task_to_run_on_fail
    
    branch1 = BranchPythonOperator(
        task_id='branch_task_1',
        provide_context=True,
        python_callable=branch_func,
        op_kwargs={'task_to_check': 'A', 'task_to_run_on_success': 'B', 'task_to_run_on_fail': 'D1'},
        trigger_rule='all_done',
    )
    
    branch2 = BranchPythonOperator(
        task_id='branch_task_2',
        provide_context=True,
        python_callable=branch_func,
        op_kwargs={'task_to_check': 'B', 'task_to_run_on_success': 'C', 'task_to_run_on_fail': 'D2'},
        trigger_rule='all_done',
    )
    

    对于依赖项:

    A >> branch_1 >> [B, D1]
    B >> branch_2 >> [C, D2]
    [C, D2, D1] >> E
    

    对于 E,您使用 trigger_rule='one_success'。 对于其他任务,您保留默认值all_success

【讨论】:

  • 如果任务 A 失败,它将为所有下游任务设置 upstream_failed,因此它永远不会执行分支任务(在您的示例中为 branch_1 或 branch_2)。
  • 没错,我刚刚更新了代码,为两个分支运算符添加了trigger_rule='all_done'
猜你喜欢
  • 2017-09-26
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2016-12-30
  • 2020-06-19
  • 2022-01-24
  • 1970-01-01
相关资源
最近更新 更多