【问题标题】:Airflow Deferrable Operator Pattern for Event-driven DAGs用于事件驱动的 DAG 的气流可延迟运算符模式
【发布时间】:2022-11-10 13:22:14
【问题描述】:

我正在寻找用于事件驱动 DAG 的模式示例,特别是那些依赖于其他 DAG 的模式。让我们从一个简单的例子开始:

dag_a -> dag_b

dag_b 取决于 dag_a。我知道在dag_a 的末尾我可以添加一个触发器来启动dag_b。然而,从抽象的角度来看,这在哲学上感觉不一致:dag_a 不需要理解或知道dag_b 的存在,但这种模式将强制要求在dag_a 上调用dag_b

让我们考虑一个稍微复杂一点的例子(原谅我糟糕的 ASCII 绘图技巧):

dag_a ------> dag_c
         /
dag_b --/

在这种情况下,如果dag_c 依赖于dag_adag_b。我知道我们可以为每个dag_adag_b 的输出设置一个传感器,但是随着可延迟运算符的出现,这似乎不是最佳实践。我想我想知道如何以异步方式设置 DAG 的 DAG。

此处的天文学家指南中介绍了事件驱动 DAG 的可延迟运算符的潜力:https://www.astronomer.io/guides/deferrable-operators,但根据上述示例,尚不清楚如何最好地应用这些运算符。

更具体地说,我正在设想一个用例,其中多个 DAG 每天运行(因此它们共享相同的运行日期),每个 DAG 的输出是某处表中的日期分区。下游 DAG 消耗上游表的分区,因此我们希望对它们进行调度,以使下游 DAG 在上游 DAG 完成之前不会尝试运行。

现在,我在下游 dag 中使用“快速且经常失败”的方法,它们在预定日期开始运行,但首先检查上游是否存在所需的数据,如果不存在则任务失败。我将这些任务设置为每隔 x 间隔重试一次,重试次数很多(例如,在 24 小时内每小时重试一次,如果仍然不存在,那么出现问题并且 DAG 失败)。这很好,因为 1)它在大多数情况下都有效,并且 2)我不相信失败的任务在重试之间继续占用一个工作槽,所以它实际上有点异步(我可能是错的)。这只是有点粗糙,所以我在想象有更好的方法。

欢迎任何有关如何将这种关系设置为更多事件驱动的战术建议,同时仍然受益于可延迟操作符的异步性质。

【问题讨论】:

  • 你有机会分享你在这个话题上的发现吗?
  • @orak 对于完全事件驱动的系统,我在这里找不到任何体面的最佳实践实践。有可能一起破解一些东西,但这似乎有点超出了 Airflow 的范式。我能想出的最好的选择就是使用可延迟的操作符来感知上游 dags 的输出。它并不完美,但效果很好。

标签: airflow directed-acyclic-graphs


【解决方案1】:

我们使用事件总线连接 DAG,DAG 的结束任务将发送事件,随后的 DAG 将根据事件类型在 Orchestrator 中触发。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-05-03
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多