Celery with Redis broker and multiple queues: all tasks are registered to each queue

Viewed 2753

I am using celery with Django and redis as the broker. I'm trying to setup two queues: default and other. My tasks are working, but the settings I have configured are not working as I am expecting them to work.

I'm having two related issues:

  1. celery tasks are not respecting the task_routes setting (see below).
  2. all of the celery tasks (no matter how they are defined) are registered to each of the two queues when they are started

Here's repo that contains all of the code. You can run the example locally with docker and docker-compose by running docker-compose up:

https://gitlab.com/verbose-equals-true/digital-ocean-docker-swarm

Here are the relevant parts of my code:

  1. the celery app definition file
  2. task definitions/declarations
  3. commands to start workers

celery app definition:

from celery import Celery
from django.conf import settings

from kombu import Exchange, Queue

CELERY_QUEUE_DEFAULT = 'default'
CELERY_QUEUE_OTHER = 'other'

app = Celery('backend')
app.conf["broker_url"] = f"redis://{settings.REDIS_SERVICE_HOST}:6379/1"
app.conf["result_backend"] = f"redis://{settings.REDIS_SERVICE_HOST}:6379/2"
app.conf["accpet_content"] = ['application/json']
app.conf["task_serializer"] = 'json'
app.conf["result_serializer"] = 'json'
app.conf["task_acks_late"] = True
app.conf["task_default_queue"] = CELERY_QUEUE_DEFAULT
app.conf["worker_send_task_events"] = True
app.conf["worker_prefetch_multiplier"] = 1
app.conf["task_queues"] = (
    Queue(
        CELERY_QUEUE_DEFAULT,
        Exchange(CELERY_QUEUE_DEFAULT),
        routing_key=CELERY_QUEUE_DEFAULT,
    ),
    Queue(
        CELERY_QUEUE_OTHER,
        Exchange(CELERY_QUEUE_OTHER),
        routing_key=CELERY_QUEUE_OTHER,
    ),
)
app.conf["task_routes"] = {
    'backend.core.tasks.debug_task': {
        'queue': 'default',
        'routing_key': 'default',
        'exchange': 'default',
    },
    'backend.core.tasks.debug_task_other': {
        'queue': 'other',
        'routing_key': 'other',
        'exchange': 'other',
    },
}
app.conf["task_default_exchange_type"] = 'direct'

app.autodiscover_tasks(lambda: settings.INSTALLED_APPS)

Task definitions (defined in a file called tasks.py in an app called core:

import time

import celery

# from backend import celery_app as app
# from celery import shared_task
from celery.task import task
from django.conf import settings

# @celery.task <-- I have seen these decorators in other example
# @app.task <-- neither of these result in the tasks being sent to the correct queue
# @shared_task
@task(queue="default", exchange="default")
def debug_task():
    time.sleep(10)
    return "Task is done."


@task(queue="other", exchange="other")
def debug_task_other():
    time.sleep(10)
    return "Task is done for other queue."

Here's how I'm starting my workers in docker-compose locally:

  celery:
    <<: *backend
    container_name: celery
    command:
      - "watchmedo"
      - "auto-restart"
      - "--directory=./"
      - "--pattern=*.py"
      - "--recursive"
      - "--"
      - "celery"
      - "worker"
      - "--app=backend.celery_app:app"
      - "-Q"
      - "default"
      - "--concurrency=1"
      - "--loglevel=INFO"
    ports: []

  celery_other:
    <<: *backend
    container_name: celery_other
    command:
      - "watchmedo"
      - "auto-restart"
      - "--directory=./"
      - "--pattern=*.py"
      - "--recursive"
      - "--"
      - "celery"
      - "worker"
      - "--app=backend.celery_app:app"
      - "-Q"
      - "other"
      - "--concurrency=1"
      - "--loglevel=INFO"
    ports: []

Here are the logs from docker-compose that show that the two tasks are both registered to each worker:

celery             |  
celery             |  -------------- celery@c8f0ed3f97df v4.4.7 (cliffs)
celery             | --- ***** ----- 
celery             | -- ******* ---- Linux-4.15.0-112-generic-x86_64-with-glibc2.2.5 2020-08-21 20:35:47
celery             | - *** --- * --- 
celery             | - ** ---------- [config]
celery             | - ** ---------- .> app:         backend:0x7fb327c2e6a0
celery             | - ** ---------- .> transport:   redis://redis:6379/1
celery             | - ** ---------- .> results:     redis://redis:6379/2
celery             | - *** --- * --- .> concurrency: 1 (prefork)
celery             | -- ******* ---- .> task events: ON
celery             | --- ***** ----- 
celery             |  -------------- [queues]
celery             |                 .> default          exchange=default(direct) key=default
celery             |                 
celery             | 
celery             | [tasks]
celery             |   . core.tasks.debug_task
celery             |   . core.tasks.debug_task_other
celery             | 
celery_other       |  
celery_other       |  -------------- celery@3dd99b1ed32e v4.4.7 (cliffs)
celery_other       | --- ***** ----- 
celery_other       | -- ******* ---- Linux-4.15.0-112-generic-x86_64-with-glibc2.2.5 2020-08-21 20:35:47
celery_other       | - *** --- * --- 
celery_other       | - ** ---------- [config]
celery_other       | - ** ---------- .> app:         backend:0x7f54fa89e6a0
celery_other       | - ** ---------- .> transport:   redis://redis:6379/1
celery_other       | - ** ---------- .> results:     redis://redis:6379/2
celery_other       | - *** --- * --- .> concurrency: 1 (prefork)
celery_other       | -- ******* ---- .> task events: ON
celery_other       | --- ***** ----- 
celery_other       |  -------------- [queues]
celery_other       |                 .> other            exchange=other(direct) key=other
celery_other       |                 
celery_other       | 
celery_other       | [tasks]
celery_other       |   . core.tasks.debug_task
celery_other       |   . core.tasks.debug_task_other

I was thinking that defining task_routes would mean that I don't have to specify the tasks's queue in the task decorator. If I don't specify the queue, the tasks are all picked up by the default worker.

If it helps, here is my Django directory structure:

tree -L 3 backend
backend
├── backend
│   ├── asgi.py
│   ├── celery_app.py <- this is where I define my celery app
│   ├── __init__.py
│   ├── settings
│   │   ├── base.py
│   │   ├── development.py
│   │   ├── __init__.py
│   │   └── production.py
│   ├── settings.py
│   ├── urls.py
│   └── wsgi.py
├── core
│   ├── admin.py
│   ├── apps.py
│   ├── __init__.py
│   ├── migrations
│   │   └── __init__.py
│   ├── models.py
│   ├── tasks.py <- this is where I define the tasks shown above
│   ├── tests.py
│   ├── urls.py
│   └── views.py
├── docker
│   ├── Dockerfile.dev
│   └── Dockerfile.prod
├── manage.py
└── requirements
    ├── base.txt
    ├── dev.txt
    └── test.txt

Here's another related SO question from a few years ago that was never answered: Celery tasks registering in multiple queues

Here's a link to the same issue that I raised on the celery/celery repo on GitHub: https://github.com/celery/celery/issues/6309

I have tried to follow the Routing Tasks page from the celery documentation to get everything setup correctly: https://docs.celeryproject.org/en/stable/userguide/routing.html

Edit:

I was mistaken about the output of celery worker process. The listed [tasks] refer to all tasks for the celery app, not the tasks that should be processed by that worker (which is what I assumed it meant). This SO post explains: How to register Celery task to specific worker?, so it makes sense that that all tasks are listed in the output of both workers in [tasks]. This solves the 2nd issue I listed above.

Also I resolved the first issue by switching to using namespaced celery settings in Django settings, I may have also had something misconfigured. The celery tasks are now routed correctly and if the route is not specified in CELERY_TASK_ROUTES then it is routed to the default queue.

1 Answers

To Solve Your First Issue, Here You Go.

In Order To Use The Settings You Provide With Celery, You Need To Use celery.update(__all__) or if you only want to update your settings, you should do celery.update(settings).

Feel free to comment if you need more help.

Related