【问题标题】:celery one broker multiple queues and workers芹菜一个经纪人多个队列和工人
【发布时间】:2017-08-15 08:01:22
【问题描述】:

我有一个名为tasks.py 的python 文件,我在其中定义了4 个单一任务。我想配置 celery 以使用 4 个队列,因为每个队列都会分配不同数量的工人。我正在阅读我应该使用 route_task 属性,但我尝试了几个选项但没有成功。

我在关注这个文档celery route_tasks docs

我的目标是运行 4 个工作人员,每个任务一个,并且不要将来自不同工作人员的任务混合在不同的队列中。这是可能的?这是个好方法吗?

如果我做错了什么,我很乐意更改我的代码以使其正常工作

这是我目前的配置

tasks.py

app = Celery('tasks', broker='pyamqp://guest@localhost//')
app.conf.task_default_queue = 'default'
app.conf.task_queues = (
    Queue('queueA',    routing_key='tasks.task_1'),
    Queue('queueB',    routing_key='tasks.task_2'),
    Queue('queueC',    routing_key='tasks.task_3'),
    Queue('queueD',    routing_key='tasks.task_4')
)


@app.task
def task_1():
    print "Task of level 1"


@app.task
def task_2():
    print "Task of level 2"


@app.task
def task_3():
    print "Task of level 3"


@app.task
def task_4():
    print "Task of level 4"

为每个队列运行 celery 一个工人

celery -A tasks worker --loglevel=debug -Q queueA --logfile=celery-A.log -n W1&
celery -A tasks worker --loglevel=debug -Q queueB --logfile=celery-B.log -n W2&
celery -A tasks worker --loglevel=debug -Q queueC --logfile=celery-C.log -n W3&
celery -A tasks worker --loglevel=debug -Q queueD --logfile=celery-D.log -n W4&

【问题讨论】:

  • 基本上我的问题是,与文档混淆,我使用的是 3.x 版本并且使用的是 4.x 的文档...epic fail

标签: python rabbitmq celery


【解决方案1】:

无需进入复杂的路由来将任务提交到不同的队列。像往常一样定义你的任务。

from celery import celery

app = Celery('tasks', broker='pyamqp://guest@localhost//')

@app.task
def task_1():
    print "Task of level 1"


@app.task
def task_2():
    print "Task of level 2"

现在在排队任务时,将任务放入适当的队列中。这是一个如何做的例子。

In [12]: from tasks import *

In [14]: result = task_1.apply_async(queue='queueA')

In [15]: result = task_2.apply_async(queue='queueB')

这会将task_1 放入名为queueAtask_2 的队列中queueB

现在您可以启动您的工人来消费它们。

celery -A tasks worker --loglevel=debug -Q queueA --logfile=celery-A.log -n W1&
celery -A tasks worker --loglevel=debug -Q queueB --logfile=celery-B.log -n W2&

【讨论】:

    【解决方案2】:

    注意:taskmessage 在答案中可以互换使用。它基本上是producer 发送到 RabbitMQ 的有效载荷

    您可以遵循 Chillar 建议的方法,也可以定义并使用 task_routes 配置将消息路由到适当的队列。这样你就不需要每次调用apply_async时都指定队列名称。

    示例:将 task1 路由到 QueueA 并将 task2 路由到 QueueB

    app = Celery('my_app')
    app.conf.update(
        task_routes={
            'task1': {'queue': 'QueueA'},
            'task2': {'queue': 'QueueB'}
        }
    )
    

    将任务发送到多个队列有点棘手。您必须声明交换,然后使用适当的routing_key 路由您的任务。您可以获取有关交换类型的更多信息here。让我们使用direct 进行说明。

    1. 创建交换

      from kombu import Exchange, Queue, binding
      exchange_for_queueA_and_B = Exchange('exchange_for_queueA_and_B', type='direct')
      
    2. 在队列上创建与该交换的绑定

      app.conf.update(
          task_queues=(
              Queue('QueueA', [
                  binding(exchange_for_queueA_and_B, routing_key='queue_a_and_b')
              ]),
              Queue('QueueB', [
                  binding(exchange_for_queueA_and_B, routing_key='queue_a_and_b')
              ])
          )
      )
      
    3. 定义task_route 发送task1到交易所

      app.conf.update(
          task_routes={
              'task1': {'exchange': 'exchange_for_queueA_and_B', 'routing_key': 'queue_a_and_b'}
          }
      )
      

    您还可以按照 Chillar 在上述答案中的建议,在您的 apply_async 方法中声明 exchangerouting_key 的这些选项。

    之后,您可以在同一台机器或不同机器上定义您的工作人员,以从这些队列中消费。

    celery -A my_app worker -n consume_from_QueueA_and_QueueB -Q QueueA,QueueB
    celery -A my_app worker -n consume_from_QueueA_only -Q QueueA
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2014-06-09
      • 2023-04-08
      • 1970-01-01
      • 2016-10-22
      • 2015-06-23
      • 2015-05-13
      • 2021-12-27
      • 2021-10-24
      相关资源
      最近更新 更多