【问题标题】:Structuring Airflow DAG with complex trigger rules使用复杂的触发规则构建 Airflow DAG
【发布时间】:2021-10-28 01:21:36
【问题描述】:

这可能是一个逻辑问题,但它让我很难过。我有以下 dag:

我有两个主要的分支事件让我绊倒:

  • 在 A 之后,只有分支 B 或 C 应该运行。
  • 在 B 之后,D 应该有选择地在 E 之前运行。

现在我已经实现了这个,以便 E 具有触发规则“none_failed”,这是为了防止在跳过 D 时跳过 E。此设置适用于所有情况,除非 C 需要运行。在这种情况下,C 和 E 同时运行(您可以通过框颜色看到 B/D 已按预期跳过,但 E 是绿色)。在阅读了触发规则文档后,我明白为什么会发生这种情况(E 的两个父任务都被跳过,因此它运行)。但是我似乎无法弄清楚如何从中获得预期的行为(例如,在运行该分支时保持 B/D/E 的当前行为,但 不要 在 C 时运行 E正在运行)。

对于其他上下文,这不是我的整个 DAG。任务 C 和 E 以 ONE_FAILED 的触发规则汇聚到它们下游的另一个任务中,但为了简单起见,我在示例中省略了这一点。任何想法如何获得预期的行为?

【问题讨论】:

  • B 或 D 失败时 E 不应该运行?
  • 没错,E不应该在1.B被跳过,或者2.B或D失败时运行。
  • 在这种情况下,E 的默认“all_success”不符合您的要求吗?
  • 不,在这种情况下,当跳过 D 时,E 不会运行。 E 的“none_failed”比“all_done”更正确,因为当 B 或 D 失败时 E 不会运行,但由于跳过的方式,仍然不能解决运行 C 时 E 运行的问题。 t 通过“all_success”或“all_failed”以外的规则级联。
  • 我更新了原始帖子以指定“none_failed”而不是“all_done”,因为这更正确但仍然不是原始问题。

标签: airflow directed-acyclic-graphs


【解决方案1】:

这可能不是最好的解决方案,但它似乎涵盖了您的所有场景。主要是我在 E 之前添加了一个虚拟任务,以控制 E 的时间并将 trigger_rule 更改为 E 为“one_success”。

“one_success”需要至少 1 个直接父级才能成功,因此对于 E,D 或 dummy 必须成功才能让 E 运行。

A = BranchPythonOperator(task_id='A', python_callable=_branch_A, dag=dag)

B = BranchPythonOperator(task_id='B', python_callable=_branch_B, dag=dag)
C = DummyOperator(task_id='C', dag=dag)

D = PythonOperator(task_id='D', python_callable=_test, dag=dag)
dummy = DummyOperator(task_id='dummy', dag=dag)
E = DummyOperator(task_id='E', trigger_rule='one_success', dag=dag)

A >> [B, C]
B >> [D, dummy] >> E

演示

【讨论】:

  • 如果我错过任何场景,请告诉我。
  • 这看起来不错,适用于预期的用例。我正在考虑在 D 和 E 之间使用带有“none_failed”的 ShortCircuitOperator 来检查 B 的状态,这将产生相同的结果,但感觉更复杂。我更喜欢你的解决方案,但我希望 Airflow 不需要这样的解决方法。
  • 我也在考虑 ShortCircuit,但在你的情况下,我无法考虑让它与 ShortCircuit 一起工作,因为你有 A >> C 没有任何 B-E 案例。无论如何,这是一个类似的案例,但不包括 A >> C 场景并且也很复杂。 stackoverflow.com/questions/51725746/…
猜你喜欢
  • 1970-01-01
  • 2018-01-16
  • 1970-01-01
  • 1970-01-01
  • 2021-12-08
  • 2020-08-14
  • 2020-05-25
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多