【发布时间】:2020-05-25 21:24:15
【问题描述】:
- 我有一组要作为 DAG 运行的作业单元(工人)
- Group1 有 10 个工作人员,每个工作人员从数据库中提取多个表。请注意,每个工作人员都映射到一个数据库实例,并且每个工作人员需要成功处理总共 100 个表才能成功将自己标记为完成
- Group1 有一个限制,即一次不应超过所有这 10 个工作人员的 5 个表。例如:
- Worker1 正在提取 2 个表
- Worker2 正在提取 2 个表
- Worker3 正在提取 1 个表
- Worker4...Worker10 需要等到 Worker1...Worker3 放弃线程
- Worker4...Worker10 可以在 step1 中的线程释放后立即拿起表
- 当每个工作人员完成所有 100 个表时,它会继续执行步骤 2,而无需等待。 Step2 没有并发限制
我应该能够创建一个满足节流的单个节点 Group1 并且还具有
- 10 个独立的工作节点,因此我可以在其中任何一个发生故障时重新启动它们
- 如果任何一个worker出现故障,我可以重新启动它而不影响其他worker。它仍然使用来自 Group1 的相同线程池,因此强制执行并发限制
- 一旦 step1 和 step2 的所有元素都完成,Group1 就会完成
- Step2 没有任何并发措施
如何在 Airflow 中为 Spring Boot Java 应用程序实现这样的层次结构? 是否可以使用 Airflow 构造设计这种 DAG,并动态地告诉 Java 应用程序一次可以提取多少个表。例如,如果除 Worker1 之外的所有工作线程都已完成,则 Worker1 现在可以使用所有 5 个可用线程,而其他所有线程都将继续执行步骤 2。
【问题讨论】:
-
您是否尝试过查看气流的连接池概念? Airflow 在 hooks 中使用连接池来控制与源的并行连接数
-
@SreenathKamath 你能指出我的文档吗?我知道 Airflow 中的任务池,但不确定您是否在这里指的是同一件事。
标签: java multithreading etl airflow airflow-scheduler