【问题标题】:How do I create celery queues on runtime so that tasks sent to that queue gets picked up by workers?如何在运行时创建 celery 队列,以便发送到该队列的任务被工作人员拾取?
【发布时间】:2014-02-10 21:10:25
【问题描述】:

我正在使用 django 1.4、celery 3.0、rabbitmq

为了描述这个问题,我在一个系统中有许多内容网络,我想要一个队列来处理与每个网络相关的任务。

但是,当系统运行时,内容是动态创建的,因此我需要动态创建队列并让现有工作人员开始处理它们。

我已经尝试通过以下方式调度任务(其中内容是 django 模型实例):

queue_name = 'content.{}'.format(content.pk)
# E.g. queue_name = content.0c3a92a4-3472-47b8-8258-2d6c8a71e3ba
add_content.apply_async(args=[content], queue=queue_name)

这将创建一个名为 content.0c3a92a4-3472-47b8-8258-2d6c8a71e3ba 的队列,创建一个名为 content.0c3a92a4-3472-47b8-8258-2d6c8a71e3ba 和路由键 content.0c3a92a4-3472-47b8-8258-2d6c8a71e3ba 的新交换,并将任务发送到该队列。

但是我从来没有看到工人接手这些任务。我当前设置的工人没有监听任何特定的队列(未使用队列名称初始化)并接手发送到默认队列就好了。我的 Celery 设置是:

BROKER_URL = "amqp://test:password@localhost:5672/vhost"
CELERY_TIMEZONE = 'UTC'
CELERY_ALWAYS_EAGER = False

from kombu import Exchange, Queue

CELERY_DEFAULT_QUEUE = 'default'
CELERY_DEFAULT_EXCHANGE = 'default'
CELERY_DEFAULT_EXCHANGE_TYPE = 'direct'
CELERY_DEFAULT_ROUTING_KEY = 'default'

CELERY_QUEUES = (
    Queue(CELERY_DEFAULT_QUEUE, Exchange(CELERY_DEFAULT_EXCHANGE),
        routing_key=CELERY_DEFAULT_ROUTING_KEY),
)

CELERY_CREATE_MISSING_QUEUES = True
CELERYD_PREFETCH_MULTIPLIER = 1

知道如何让工作人员接手发送到这个新创建队列的任务吗?

【问题讨论】:

  • 为什么不只使用一个队列并将content.pk 作为参数传递呢?创建新队列还有什么好处?
  • 可能的附加好处是:如果需要,可以为接收大量流量的内容网络启动专门的工作人员。也用于统计和日志封装等。

标签: django queue celery amqp worker


【解决方案1】:

我们可以动态添加队列并将工作器附加到它们。

from celery import current_app as app
from task import celeryconfig #your celeryconfig module

动态定义任务并将其路由到队列中

from task import process_data
process_data.apply_async(args,kwargs={}, queue='queue-name')
reply = app.control.add_consumer('queue_name', destination = ('your-worker-name',), reply = True)

您必须将队列名称保存在像 redis 这样的持久数据存储中,以便在重新启动时记住它。

redis.sadd('CELERY_QUEUES','queue_name')

celeryconfig.py 也使用相同的方法来记住队列名称

CELERY_QUEUES = {
    'celery-1': {
        'binding_key': 'celery-1'
    },
    'gateway-1': {
        'binding_key': 'gateway-1'
    },
    'gateway-2': {
        'binding_key': 'gateway-2'
    }
}
for queue in redis.smembers('CELERY_QUEUES'):
        CELERY_QUEUES[queue] = dict(binding_key=queue)

【讨论】:

    【解决方案2】:

    您需要告诉工作人员开始使用新队列。相关文档为here

    从命令行:

    $ celery control add_consumer content.0c3a92a4-3472-47b8-8258-2d6c8a71e3ba
    

    或者从python内部:

    >>> app.control.add_consumer('content.0c3a92a4-3472-47b8-8258-2d6c8a71e3ba', reply=True)
    

    两种形式都接受目标参数,因此如果需要,您可以只告诉个别工作人员有关新队列的信息。

    【讨论】:

    • 尽管正如@EWit 所建议的那样,除非您有充分的理由将作业分配给不同的工作人员,否则我会坚持使用默认队列机制。 add_consumer 调用不是“粘性” - 如果您重新启动工作人员,您将必须重新应用队列订阅。
    • 谢谢加文。我现在正在检查 add_consumer 方法。如果我可以让它工作,会告诉你。
    • 太棒了!非常感谢加文。重新创建队列不是问题,因为它们总是在必要时动态创建。
    • 不客气。我只是想从我之前的评论中强调,如果你重新启动一个工作器(例如,当有代码更改或服务器重新启动时)它只会开始使用默认队列。每次启动工作人员时,您都必须执行 add_consumer。
    猜你喜欢
    • 2013-02-18
    • 2020-01-28
    • 2018-07-18
    • 1970-01-01
    • 2019-02-22
    • 2015-06-16
    • 2012-06-30
    • 2013-06-02
    相关资源
    最近更新 更多