【问题标题】:How to execute plain Celery tasks on Airflow workers如何在 Airflow 工作人员上执行普通的 Celery 任务
【发布时间】:2019-10-05 09:30:10
【问题描述】:

我目前使用 CeleryExecutor 作为后端来设置 Airflow 并正常工作,以提供水平缩放。这非常有效,尤其是当工作节点位于 EC2 上的自动缩放组中时。

除了 Airflow,我还使用普通的 Celery 来处理来自 Flask/Python 的简单异步任务(不需要整个管道)。到目前为止,这些普通的 Celery 任务的数量非常少,我只是在与 Flask 相同的机器上运行普通的 Celery 工作程序。现在需要在系统中运行大量的普通 Celery 任务,所以我也需要扩展我的普通 Celery。

这样做的一种方法是在 Airflow 工作服务器上运行普通的 Celery 工作服务(以从自动缩放等中受益),但这似乎不是一个优雅的解决方案,因为它创建了两个不同的“类型”的芹菜工人在同一台机器上。我的问题是,我是否可以将某些配置设置组合传递给我的普通 Celery 应用程序,这将导致 @celery.task 装饰功能作为普通 Celery 任务直接在我的 Airflow 工作集群上执行,完全绕过 Airflow 中间件。

感谢您的帮助。

【问题讨论】:

    标签: python flask celery airflow


    【解决方案1】:

    如果我没记错的话,申请是airflow.executors.celery_executor.app。尝试celery -A airflow.executors.celery_executor.app inspect active 以在您当前的 Airflow 基础架构中进行测试。但是,我建议您不要这样做,因为您的 Celery 任务可能会影响 Airflow DAG 的执行,并且可能会影响 SLA。

    我们在我工作的公司所做的正是您所建议的 - 我们维护一个大型 Celery 集群,有时我们会将一些 Airflow 任务的执行卸载到我们的 Celery 集群,具体取决于用例。当我们的 Airflow DAG 中的任务实际上触发了数以万计的小作业时,这特别方便。我们的 Celery 集群在繁忙的一天运行 800 万个任务。

    【讨论】:

    • 感谢您的回复。如果按照第一种方式完成,您能否详细说明 Celery 任务如何影响 Airflow DAG?是因为它们给服务器增加了额外的负载,还是因为它们通过了 Airflow 基础设施?另外,当你说你们按照我的建议做的时候,你的意思是你在一个节点上运行两种风格的 Celery worker,还是每个节点都有不同类型的节点?
    • 因为您可能会占用所有可用的工作进程。而那些重要的会等待并且可能会错过 SLA。
    猜你喜欢
    • 2021-08-15
    • 2014-04-13
    • 2012-11-30
    • 2016-02-14
    • 2018-10-07
    • 1970-01-01
    • 1970-01-01
    • 2021-08-14
    • 2016-11-11
    相关资源
    最近更新 更多