【问题标题】:How to run Airflow DAG for specific number of times?如何运行 Airflow DAG 特定次数?
【发布时间】:2018-12-14 06:57:13
【问题描述】:

如何运行气流 dag 指定次数?

我尝试使用 TriggerDagRunOperator,该运算符对我有用。 在可调用函数中,我们可以检查状态并决定是否继续。

但是需要维护当前的计数和状态。

使用上述方法,我可以重复 DAG 'run'。

需要专家意见,有没有其他深刻的方法来运行 Airflow DAG X 次? 谢谢。

【问题讨论】:

  • 试着解释一下你想要实现的工作需要运行特定次数。
  • 假设单个 DAG 有 5 个任务。我想运行这个 DAG 10 次。假设单次运行有时需要 2 个小时以上。只是我不能根据时间安排。因此希望根据我将在 DAG 配置中指定的数字运行 DAG。

标签: airflow google-cloud-composer


【解决方案1】:

恐怕 Airflow 完全是基于时间的调度。
您可以将计划设置为None,然后使用API to trigger runs,但您将在外部执行此操作,从而维护确定何时以及为何触发外部的计数和状态。

当您说您的 DAG 可能有 5 个任务要运行 10 次并且运行需要 2 个小时并且您无法根据时间安排它时,这会令人困惑。我们不知道 2 小时对您有什么意义,或者为什么它必须是 10 次运行,也不知道为什么您不能安排它每天运行一次这 5 个任务。使用简单的每日计划,它将每天在大约相同的时间运行一次,并且在任何一天花费的时间都不会超过 2 小时。对吧?

您可以将 start_date 设置为 11 天前(虽然是固定日期,但不要动态设置),将 end_date 设置为今天(也是固定的),然后添加每日 schedule_interval 和 @ 987654327@ of 1,您将获得 10 次运行,它会在相应地更改 execution_date 时将它们背靠背运行而不会重叠,然后停止。或者,您可以将 airflow backfillNone 计划的 DAG 和一系列执行日期时间一起使用。

您的意思是您希望它每 2 小时连续运行一次,但有时它会运行更长时间并且您不希望它重叠运行?好吧,您绝对可以安排它每 2 小时运行一次 (0 0/2 * * *) 并将 max_active_runs 设置为 1,这样如果前一次运行尚未完成,下一次运行将等待,然后在前一次运行完成时启动.请参阅https://airflow.apache.org/faq.html#why-isn-t-my-task-getting-scheduled 中的最后一个项目符号。

如果您希望您的 DAG 恰好每 2 小时运行一次 [给予或接受一些调度程序延迟,是的,这是一件事] 并让之前的运行继续进行,这主要是默认行为,但您可以添加 @987654333 @ 到一些本身不应该同时运行的重要任务(例如创建、插入或删除临时表),或使用具有单个插槽的池。

如果您的下一个计划已准备好开始,则没有任何功能可以终止先前的运行。如果之前的运行尚未完成,可能会跳过当前运行,但我忘记了具体是如何完成的。

这基本上是您的大部分选择。您也可以为计划外的 DAG 创建手动 dag_runs;一次创建 10 个(使用 UI 或 CLI 而不是 API,但 API 可能更容易)。

这些建议是否能解决您的顾虑?由于不清楚您为什么需要固定的运行次数、频率或使用什么计划和条件,因此很难提供具体的建议。

【讨论】:

  • 非常有帮助的答案,从第四段中找到了一些线索。非常感谢你。基本上,我想交叉检查是否有任何设置或配置(编码或基于 json)依次运行 DAG。例如,只是为了类比,考虑我们要在“for-loop”中运行 DAG,其中“for-loop”是使用计数器变量控制的。如果这样的功能可用,那么我不会担心执行单次迭代所需的时间。因此,用户可以手动运行并通过气流变量在外部控制迭代次数。
  • @Omkara 从您的评论来看,您可能想尝试在 BranchOperator 中结束您的 DAG,这将分支到 Dummy END 任务或 TriggerDagRunOperator 在其自己的 DAG id 和这会减少 Airflow 变量或其他一些外部数据源(DB、http get/put/post、S3/GCP 路径中的值等)以确定分支路径。
  • 是的,正如我之前在帖子中提到的,这是我第一次尝试时所做的。但是对我来说看起来有点笨拙。无论如何。
【解决方案2】:
  • Airflow 本身不支持此功能
  • 但是通过利用元数据库,我们可以自己制作这个功能

我们可以写一个自定义操作符/python 操作符

  • 在运行实际计算之前,检查元数据库中是否已经存在为任务(TaskInstance 表)运行的“n”个。 (请参阅task_command.py 寻求帮助)
  • 如果他们这样做,请跳过任务(raise AirflowSkipExceptionreference

这篇优秀的文章可以启发灵感:Use apache airflow to run task exactly once


注意

这种方法的缺点是它假定任务的历史运行 (TaskInstances) 将永远保留(并且正确)

  • 但在实践中,我经常发现 task_instances 缺失(我们将 catchup 设置为 False
  • 此外,在大型 Airflow 部署中,可能需要设置元数据库的例行 cleanup,这将使这种方法变得不可能

【讨论】:

    猜你喜欢
    • 2021-06-19
    • 2022-01-18
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-06-07
    相关资源
    最近更新 更多