【问题标题】:Implementing branching in Airflow在 Airflow 中实现分支
【发布时间】:2020-11-17 16:46:26
【问题描述】:

我正在尝试向我的气流数据添加警报。 dags 有多个任务,有些多达 15 个。 我想执行一个 bash 脚本(所有 dags 的通用脚本),以防任何任务在任何时候失败。 例如,dag 具有任务 T1 到 T5,如 T1 >> T2 >> T3 >> T4 >> T5。 如果其中任何一个失败,我想触发任务 A(代表警报)。

对于任何可以帮助我处理任务层次结构的人都会很有帮助。

【问题讨论】:

    标签: airflow


    【解决方案1】:

    IMO 有两种选择。失败回调和触发规则

    成功/失败回调

    Airflow 任务实例具有在失败或成功的情况下如何处理的概念。这些是在任务达到特定状态时运行的回调......这是您的选择:

    ...
        on_failure_callback=None,
        on_success_callback=None,
        on_retry_callback=None
    ...
    

    触发规则

    Airflow 任务实例通过default being ALL_SUCCESS 了解其upstream to trigger on 的状态。这意味着您的主分支可以保持原样。您可以使用 T1 中的 A 分支到您想要的位置:

    from airflow.utils.trigger_rule import TriggerRule
    
    T1 >> DummyOperator(
        dag=dag,
        task_id="task_a",
        trigger_rule=TriggerRule.ALL_FAILED
    )
    

    或者,您可以构建您的分支并将 A 包含为:

    from airflow.utils.trigger_rule import TriggerRule
    
    [T1, T2, T3, ...] >> DummyOperator(
        dag=dag,
        task_id="task_a",
        trigger_rule=TriggerRule.ONE_FAILED
    )
    

    【讨论】:

      猜你喜欢
      • 2020-02-11
      • 1970-01-01
      • 2022-08-11
      • 1970-01-01
      • 2019-06-08
      • 2010-10-04
      • 2022-11-22
      • 1970-01-01
      相关资源
      最近更新 更多