路由 celery 任务

gon*_*zor 0 python django celery

当尝试创建两个单独的专用工作人员时,我无法将任务发送到芹菜。我已经阅读了文档和这个问题,但它并没有改善我的情况。

我的配置如下:

CELERY_RESULT_BACKEND = 'django-db'
CELERY_BROKER_URL = f'redis://{env("REDIS_HOST")}:{env("REDIS_PORT")}/{env("REDIS_CELERY_DB")}'
CELERY_DEFAULT_QUEUE = 'default'
CELERY_DEFAULT_EXCHANGE_TYPE = 'topic'
CELERY_DEFAULT_ROUTING_KEY = 'default'
CELERY_QUEUES = (
    Queue('default', Exchange('default'), routing_key='default'),
    Queue('media', Exchange('media'), routing_key='media'),
)
CELERY_ROUTES = {
    'books.tasks.resize_book_photo': {
        'queue': 'media',
        'routing_key': 'media',
    },
}
Run Code Online (Sandbox Code Playgroud)

任务在文件中按以下方式定义tasks.py:

import logging
import time

from celery import shared_task


from books.models import Author, Book
from books.commands import resize_book_photo as resize_book_photo_command


logger = logging.getLogger(__name__)


@shared_task
def list_test_books_per_author():
    time.sleep(5)
    queryset = Author.objects.all()
    for author in queryset:
        for book in author.testing_books:
            logger.info(book.title)


@shared_task
def resize_book_photo(book_id: int):
    resize_book_photo_command(Book.objects.get(id=book_id))
Run Code Online (Sandbox Code Playgroud)

他们被称为使用apply_async:

list_test_books_per_author.apply_async()
resize_book_photo.apply_async((book.id,))
Run Code Online (Sandbox Code Playgroud)

当我运行芹菜花时,我看到队列中没有出现任何任务。 芹菜花面板显示两名没有任务的活跃工人。

工人们开始使用:

celery -A blacksheep worker -l info --autoscale=10,1 -Q media --host=media@%h
celery -A blacksheep worker -l info --autoscale=10,1 -Q default --host=default@%h
Run Code Online (Sandbox Code Playgroud)

我能做的是通过使用redis-cli和127.0.0.1:6379> LRANGE celery 1 100命令确认它们最终在celery密钥下(这是芹菜的默认密钥)。似乎没有工人消费。

编辑仔细查看这部分文档后,我发现我的命名是错误的。将设置更改为后:

CELERY_RESULT_BACKEND = 'django-db'
CELERY_BROKER_URL = f'redis://{env("REDIS_HOST")}:{env("REDIS_PORT")}/{env("REDIS_CELERY_DB")}'
CELERY_TASK_DEFAULT_QUEUE = 'default'
# CELERY_DEFAULT_EXCHANGE_TYPE = 'topic'
CELERY_TASK_DEFAULT_ROUTING_KEY = 'default'
CELERY_QUEUES = (
    Queue('default', Exchange('default'), routing_key='default'),
    Queue('media', Exchange('media'), routing_key='media'),
)
CELERY_ROUTES = {
    'books.tasks.resize_book_photo': {
        'queue': 'media',
        'routing_key': 'media',
    },
}
Run Code Online (Sandbox Code Playgroud)

情况得到改善:任务是从default队列中消耗的,但我想要进入media队列的任务也进入了default.

EDIT2我试图通过将其调用更改为 来显式告诉某些任务转到其他队列resize_book_photo.apply_async((book.id,), queue='media')。任务已正确分派到正确的队列并被消耗。但是,我希望自动处理此问题,这样我就不必在每次调用时都定义队列apply_async

小智 7

尝试CELERY_TASK_ROUTES代替CELERY_ROUTES. 这最近对我来说与 django 集成很有用。

解释埋在这个评论中:How to routetasks to differentqueues with Celery and Django