【问题标题】:Django & Celery — Routing problemsDjango & Celery——路由问题
【发布时间】:2012-05-22 17:49:25
【问题描述】:

我正在使用 Django 和 Celery,并且正在尝试设置到多个队列的路由。当我指定任务的 routing_keyexchange(在任务装饰器中或使用 apply_async())时,任务不会添加到代理(即 Kombu 连接到我的 MySQL 数据库)。

如果我在任务装饰器中指定队列名称(这意味着路由键被忽略),任务工作正常。路由/交换设置似乎有问题。

知道可能是什么问题吗?

设置如下:

settings.py

INSTALLED_APPS = (
    ...
    'kombu.transport.django',
    'djcelery',
)
BROKER_BACKEND = 'django'
CELERY_DEFAULT_QUEUE = 'default'
CELERY_DEFAULT_EXCHANGE = "tasks"
CELERY_DEFAULT_EXCHANGE_TYPE = "topic"
CELERY_DEFAULT_ROUTING_KEY = "task.default"
CELERY_QUEUES = {
    'default': {
        'binding_key':'task.#',
    },
    'i_tasks': {
        'binding_key':'important_task.#',
    },
}

tasks.py

from celery.task import task

@task(routing_key='important_task.update')
def my_important_task():
    try:
        ...
    except Exception as exc:
        my_important_task.retry(exc=exc)

启动任务:

from tasks import my_important_task
my_important_task.delay()

【问题讨论】:

  • 如何传递routing_key?使用 async_apply?
  • 我正在使用delay() 方法,它只是apply_async() 的快捷方式。我正在尝试使用任务方法(通过装饰器)而不是在调用它时保留routing_key 规范。我尝试使用 apply_async() 传递密钥,但我遇到了同样的问题。
  • delay 不接受 routing_key 关键字。它是 apply_async 的简化版本,但它们并不相同。
  • 我在启动任务时没有传递任何路由信息,我在任务装饰器中指定了它,并带有方法定义。请参阅上面的代码以查看我的设置。
  • 我是否认为我可以在任务装饰器中指定路由/交换信息并且在调用时应该尊重它?

标签: django celery django-celery kombu


【解决方案1】:

您正在使用 Django ORM 作为代理,这意味着声明仅存储在内存中 (请参阅http://readthedocs.org/docs/kombu/en/latest/introduction.html#transport-comparison 上的交通比较表,很难找到)

因此,当您使用 routing_key important_task.update 应用此任务时,它将无法 路由它,因为它还没有声明队列。

如果你这样做,它会起作用:

@task(queue="i_tasks", routing_key="important_tasks.update")
def important_task():
    print("IMPORTANT")

但是使用自动路由功能对您来说会简单得多, 因为这里没有任何内容表明您需要使用“主题”交换, 要使用自动路由,只需删除设置

  • CELERY_DEFAULT_QUEUE,
  • CELERY_DEFAULT_EXCHANGE,
  • CELERY_DEFAULT_EXCHANGE_TYPE
  • CELERY_DEFAULT_ROUTING_KEY
  • CELERY_QUEUES

然后像这样声明你的任务:

@task(queue="important")
def important_task():
    return "IMPORTANT"

然后启动一个从该队列消费的工作人员:

$ python manage.py celeryd -l info -Q important

或者从默认 (celery) 队列和 important 队列中消费:

$ python manage.py celeryd -l info -Q celery,important

另一个好的做法是不要将队列名称硬编码到 任务并改用CELERY_ROUTES

@task
def important_task():
    return "DEFAULT"

然后在您的设置中:

CELERY_ROUTES = {"myapp.tasks.important_task": {"queue": "important"}}

如果您仍然坚持使用主题交换,那么您可以 添加此路由器以在第一时间自动声明所有队列 发送任务:

class PredeclareRouter(object):
    setup = False

    def route_for_task(self, *args, **kwargs):
        if self.setup:
            return
        self.setup = True
        from celery import current_app, VERSION as celery_version
        # will not connect anywhere when using the Django transport
        # because declarations happen in memory.
        with current_app.broker_connection() as conn:
            queues = current_app.amqp.queues
            channel = conn.default_channel
            if celery_version >= (2, 6):
                for queue in queues.itervalues():
                    queue(channel).declare()
            else:
                from kombu.common import entry_to_queue
                for name, opts in queues.iteritems():
                    entry_to_queue(name, **opts)(channel).declare()
CELERY_ROUTES = (PredeclareRouter(), )

【讨论】:

  • Celery 3 中的队列声明和交换问题是否已解决?我在设置中使用了新的CELERY_QUEUES = (Queue(...), ...),这是否意味着队列声明正确?
  • 注意:在 Celery 4.0 以后,CELERY_ROUTES 已被 CELERY_TASK_ROUTES 替换。可能会节省某人的时间。
猜你喜欢
  • 2013-09-11
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2020-10-09
  • 2020-04-06
  • 1970-01-01
相关资源
最近更新 更多