【问题标题】:Airflow - How to have a worker take all dag run tasks?气流 - 如何让工作人员完成所有 dag 运行任务?
【发布时间】:2019-07-21 04:27:39
【问题描述】:

我目前正在使用 Airflow 和 Celery 处理文件。工作人员需要下载文件、处理它们并在之后重新上传它们。我的 DAG 只有一名工人就可以了。但是当我添加一个时,事情就变得复杂了。

工人在有空时接受任务。 Worker1 可以接任务“处理下载的文件”,但 Worker2 接任务“下载文件”,所以任务失败,因为它无法处理不存在的文件。

有没有办法向工作人员(或调度程序)指定 DAG 必须仅在一个工作人员上运行?我知道队列。但我已经在使用它们了。

【问题讨论】:

  • 避免此问题的最简单方法是将下载和处理合并到一个 task 中。顺便说一句,queues 有什么问题?
  • 感谢您的回答,但实际上有 3 个以上的任务,因此合并可能不是一个好主意。队列的问题是我实际上使用它们,并且我希望我的气流安装是可扩展的,而不是每次添加或删除一个工人时都更改我的 DAG 文件......

标签: celery airflow


【解决方案1】:

在这种情况下,您可以使用 Airflow 变量来保存所有工作节点名称。 例如:

  • 变量:worker_list
  • 值:boxA, boxB, boxC

在运行 Airflow 工作器时,您可以指定多个作业队列。例如:airflow worker job_queue1,job_queue2 对于您的情况,我将运行 airflow worker af_<hostname>

在您的 DAG 代码中,只需要获取该 worker_list Airflow 变量,随机选择一个框,然后将所有作业排队到 af_<random_selected_box> queue

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2018-09-26
    • 1970-01-01
    • 2012-08-11
    • 1970-01-01
    • 1970-01-01
    • 2019-12-06
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多