【发布时间】: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_a 和dag_b。我知道我们可以为每个dag_a 和dag_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