【发布时间】:2019-03-25 14:08:09
【问题描述】:
我有大约一百万个使用相同 python 函数的 Airflow 任务。每个都需要以不同的开始日期和参数运行。
之前我向question 询问了如何在一个 DAG 下运行两个这样的任务。但是,当任务变得很多时,那里的答案是不可扩展的。 (见链接和注释)
问题
如何在 Airflow 上以可扩展的方式运行一百万个(或任何大量)或任务,其中每个任务源自相同的 python 函数但具有不同的开始日期和不同的参数?
备注
这些任务不需要在 PythonOperator 上运行(因为它们源于 python 函数)。实际上,它们很可能会在 Kubernetes 集群上以分布式方式运行(因此使用 KubernetesExecutor 或 KubernetesPodOperator)。无论哪种方式,DAG 贡献背后的架构问题仍然存在。)
解决方案思路
我正在考虑的一个解决方案是,在一个 DAG 下,动态构建所有任务并在执行的 python 函数中传递不同的开始日期。在外部 Airflow 每天都会执行每个任务,但在函数内部,如果 execution_date 早于 start_date,则函数将只是 return 0。
【问题讨论】:
-
你能提供更多细节吗?这听起来像是很多工作。你准备好向它投掷一支机器大军,让它在宇宙热寂之前完成吗?
-
好的,请告诉我要添加什么样的信息?我现在把它限制在一个明确的问题上。
-
我有点困惑。我曾经在一家有类似工作流程安排的公司工作。我无法想象任务的数量是数千,更不用说数百万了,所以我认为我在这里遗漏了一些东西(至少三个数量级)。
-
这里有一个例子:假设我有一百万用户。他们每个人都在不同的日期加入我的网络(因此有不同的开始日期)。每个人的活动都保存在每日 .json 文件中。如果我想下载这个日期来工作,我需要为每个用户设置一个任务。它们都有不同的开始日期,我用来下载的函数会有不同的参数(例如用户名)。您的评论表明我可能以错误的方式思考这个问题。
-
你想错了。气流可以用于数百万个动态任务,但它不应该。气流 DAG 应该是相当稳定的。我建议你使用其他工具来解决这个问题。例如,您仍然可以使用 Airflow 来处理整个用户群,并稍后在您的 ETL 流程中使用此信息。
标签: airflow