Message queues and background jobs: Celery, Redis, Kafka

Move slow work off the request path: queues vs brokers vs streams, Celery in Django, retries and idempotency, dead letters, separate queues, and when Kafka is worth it.

9 min read
On this page 8 sections
  1. What belongs in a background job
  2. Queues, brokers and streams
  3. Celery in a Django app
  4. Retries, idempotency and dead letters
  5. Priorities and separate queues
  6. Kafka vs RabbitMQ vs SQS
  7. Key takeaways
  8. Frequently asked questions

A message queue holds units of work, called messages, until a separate worker process picks them up, so slow tasks run in the background instead of inside a user's request. In a Django learning platform, that means a student's "Submit test" returns in milliseconds while scoring, result emails, certificates and video transcoding happen afterwards. Celery with Redis or RabbitMQ covers most needs; Kafka earns its extra complexity only when many systems need to read the same stream of events.

What belongs in a background job

Move work off the request path when the user doesn't need its result in the response, when it's slow, when it depends on a third party that can be slow or down, or when one action fans out into many:

TaskWhy it goes to a queue
Scoring a submitted test and updating ranksThe answers must be saved immediately; the score can follow a few seconds later
SMS, WhatsApp and email notificationsGateways can be slow or rate-limited, and failures need retries
Notifying 20,000 students that a lecture is liveOne action fans out into thousands of messages
Transcoding a lecture, compressing uploaded PDFs, generating certificatesMinutes of compute per item
Reports, exports and analytics rollupsHeavy queries that shouldn't compete with students at peak

Keep in the request anything the user must see confirmed: the answers themselves are saved in the database before "Submit" returns, and a payment is verified before access is granted. The queue handles what comes after.

Queues, brokers and streams

These words get used interchangeably, but they describe different things:

  • A queue is an ordered buffer of messages; each message is normally processed by one consumer and then removed.

  • A message broker is the server software that accepts messages from producers, stores them in queues and delivers them to consumers. RabbitMQ is a full broker with routing rules; Redis is often used as a simple one.

  • A stream or log keeps messages after they are read. As Kafka's introduction puts it, events are not deleted after consumption; each topic keeps them for a configured period, and different consumer groups read at their own pace.

So "message queue vs message broker" isn't a choice between two products: the broker is the system, and the queues live inside it.

Celery in a Django app

Celery is a mature, widely used task queue for Python and Django. The standard setup is a small module that creates the app and loads its settings from Django:

# proj/celery.py
import os
from celery import Celery

os.environ.setdefault("DJANGO_SETTINGS_MODULE", "proj.settings")
app = Celery("proj")
app.config_from_object("django.conf:settings", namespace="CELERY")
app.autodiscover_tasks()

Tasks are plain functions with a decorator. This one retries only on a gateway error, backing off exponentially:

from celery import shared_task

@shared_task(bind=True, autoretry_for=(SmsGatewayError,),
             retry_backoff=True, max_retries=5, acks_late=True)
def send_result_sms(self, attempt_id):
    attempt = Attempt.objects.select_related("student").get(pk=attempt_id)
    if attempt.result_sms_sent_at:      # already sent: skip on a retry
        return
    sms_gateway.send(attempt.student.phone, attempt.result_text())
    Attempt.objects.filter(pk=attempt_id).update(result_sms_sent_at=now())

Three habits prevent most bugs:

  • Pass IDs, not model objects. Celery's documentation advises refetching data inside the task, because the object may have changed since the task was queued.

  • Enqueue after the transaction commits. Otherwise a worker can pick up the task before the row it needs exists. Celery 5.4 added delay_on_commit() for Django; before that, wrap the call in transaction.on_commit().

  • Choose the broker deliberately. With Redis, a task that isn't acknowledged within the visibility timeout, one hour by default, is redelivered to another worker, and Celery's Redis notes warn that long ETA or countdown tasks can run repeatedly as a result. Use a dedicated Redis that won't evict keys, or RabbitMQ for long-running work.

Django 6.0 added a built-in Tasks framework for defining and queuing background tasks. It standardises the API, but it doesn't run the work: production still needs a backend and worker process, from a third-party package, to execute queued tasks.

Retries, idempotency and dead letters

Queues deliver messages at least once, not exactly once. A worker can crash after doing the work but before acknowledging it; with acks_late enabled, Celery will then run the task again, and its documentation says plainly that tasks must be idempotent. Amazon SQS standard queues make the same at-least-once promise. Design for duplicates:

  • Make tasks idempotent. Check-and-mark as in the example above, use unique constraints (one certificate per student per course), and prefer "set status to X" over "increment". A crash between sending and recording can still cause a repeat, so where a provider accepts an idempotency key, pass one, such as the attempt ID.

  • Retry only transient errors. A timeout from an SMS gateway is worth retrying; an invalid phone number isn't. Use exponential backoff with jitter; Celery's retry_backoff caps the delay at 600 seconds and adds jitter by default.

  • Park what keeps failing. A dead-letter queue holds messages that failed too often, for a person to inspect. RabbitMQ can dead-letter messages that are rejected, expire or exceed a quorum queue's delivery limit; SQS moves a message after maxReceiveCount receives. Celery on Redis has no built-in equivalent, so record final failures in a table and alert on it.

A worked example, with illustrative numbers: on results day, 20,000 students each need an SMS, and your gateway accepts 100 messages a second. Draining the queue takes about 200 seconds, which is fine, as long as you respect the limit. Celery's rate_limit option applies per worker instance, not globally, so with four workers each would need 25 per second. The simpler route to a global limit is a dedicated queue with a known number of workers.

Priorities and separate queues

One queue for everything means a backlog of report exports can delay OTPs. Route tasks to queues by urgency and duration, and run separate worker pools for each:

# settings.py
CELERY_TASK_ROUTES = {
    "accounts.tasks.send_otp": {"queue": "critical"},
    "exams.tasks.score_attempt": {"queue": "default"},
    "videos.tasks.transcode_lecture": {"queue": "video"},
    "reports.tasks.*": {"queue": "bulk"},
}

# run separate pools, for example:
#   celery -A proj worker -Q critical --concurrency=8
#   celery -A proj worker -Q video --concurrency=1

Give the video queue its own machines, late acknowledgement and a prefetch multiplier of 1. Celery's default multiplier of 4 lets one worker reserve several long jobs while others sit idle, and its optimisation guide recommends separate workers for long and short tasks. Watch queue length and the age of the oldest message for every queue, and scale workers on those numbers; our guide to autoscaling workers with HPA and KEDA shows how. For what the transcoding workers actually do, see our guide to video transcoding pipelines.

Kafka vs RabbitMQ vs SQS

AspectRabbitMQAmazon SQSApache KafkaRedis as a broker
ModelBroker with exchanges, routing and queuesManaged queue servicePartitioned, replicated logIn-memory data store used as a queue
After a message is processedRemoved once acknowledgedDeleted by the consumerKept for the topic's retention period; can be replayedRemoved
OrderingPer queueBest-effort (standard); strict per message group (FIFO)Strict within a partitionPer list or stream
OperationsA cluster to runNone; pay per requestThe heaviest to run, unless managedOften already in your stack
Good forTask queues with routing and dead-letteringTask queues on AWS without running a brokerEvent streams read by many consumersSimple, fast task queues

A few specifics shape the choice. SQS hides a received message for a visibility timeout of 30 seconds by default, extendable up to 12 hours, and FIFO queues handle up to 3,000 messages a second per API action with batching, or more in high-throughput mode. In Kafka's classic consumer groups, each partition is read by exactly one consumer in the group, so the partition count caps parallelism; recent versions also add share groups, in which several consumers can read the same partition.

Kafka is worth it when the same events feed several independent systems. Playback heartbeats from thousands of students, for example, might drive live analytics, watch-time billing, recommendations and an audit trail, each reading the same stream at its own pace and able to replay history. For "send this email once", it's the wrong tool; a task queue is simpler and does exactly that. For how queues fit into a platform built for peak load, read how to handle 100,000 concurrent users, and for Redis as a cache rather than a broker, see our guide to Redis caching.

Key takeaways

  • Queue anything slow, flaky or fan-out that the user doesn't need in the response; keep saving answers and verifying payments in the request.

  • Delivery is at least once: make every task idempotent and retry only transient errors, with backoff.

  • Enqueue after the database transaction commits, and pass IDs rather than objects.

  • Separate queues and worker pools by urgency and duration; alert on queue age.

  • Use Celery with Redis or RabbitMQ for tasks; bring in Kafka when many consumers need the same event stream.

Frequently asked questions

What is a message queue used for?

A message queue is used to hand work from one part of a system to another without making either wait. Typical uses are sending notifications, processing uploads, scoring tests, generating reports and absorbing traffic spikes: the web server adds a message and responds at once, and workers process the backlog at a steady pace. It also lets producers and consumers be deployed and scaled independently.

What is a message queue in system design?

In system design, a message queue is the component that decouples producers from consumers. It buffers bursts so downstream services aren't overwhelmed, lets slow work happen asynchronously, and keeps messages safe while a consumer is down. Designs that use one must also handle at-least-once delivery, ordering, retries, dead-letter queues and monitoring of queue depth and message age.

What is a message queue and how does it work?

A message queue is a buffer, usually run by a broker, that stores messages until a consumer processes them. A producer sends a message; the broker stores it; a consumer receives it, does the work and acknowledges it, and the broker deletes it. If the consumer crashes before acknowledging, the message becomes available again and another consumer processes it, which is why tasks should be safe to run twice.

Is Redis a message queue?

Redis is an in-memory data store, not a dedicated message broker, but its lists and streams make it a popular lightweight queue, and Celery supports it as a broker. It is fast and often already in your stack. The trade-offs are durability, which depends on Redis's persistence settings, and Celery's visibility-timeout behaviour for long tasks, so give the broker its own Redis instance.

Share this article

Looking for something else?

Talk to Us