Available for day contractsFrom 21st September I have availability for day and half day contracts. Please contact for more information.

Contact →
mikepreston.org

RabbitMQ Without Celery: Queue Patterns for Python People Who Want to Understand Their Own System

A calm mid-century postal sorting room with brass pneumatic chutes and tidy pigeonholes that one operator understands completely, while an ornate over-engineered machine sits switched off in the corner — 1960s gouache.

There is a moment in most incidents involving background work where someone says, with more hope than confidence, "I think the task ran twice." And then nobody in the room can say for certain whether it did, or why, or what would have happened if it had. The job is a Celery task. The broker is RabbitMQ. And the gap between those two facts — the abstraction that was supposed to make life easier — is precisely where the understanding went to die.

This is not a complaint about Celery's code quality, which is fine. It is a complaint about what the abstraction costs you at exactly the wrong time. Celery's pitch is that you should not have to think about the broker: you write a function, decorate it, call .delay(), and a worker somewhere picks it up. For a great many teams that is all they need, right up until the afternoon it isn't, at which point they discover they have been running a distributed system whose delivery semantics they have never once had to articulate. The broker was there the whole time. They just never had to look at it.

I run RabbitMQ without Celery, deliberately, and have done for years. Not because the queue patterns are exotic — they are not — but because the patterns are the system. Once you can name an exchange, a queue, a binding, an ack, and a prefetch count, and say what each does when a consumer falls over mid-message, you understand your own delivery guarantees. That understanding is worth more than the convenience you gave up to get it.

Work queues from first principles

Strip away the framework and a RabbitMQ work queue is almost embarrassingly simple. A producer publishes a message to an exchange. The exchange routes it, according to bindings, into one or more queues. A consumer subscribes to a queue and receives messages. That is the whole model, and the three nouns — exchange, queue, binding — are the things Celery papers over and you should not.

The exchange is the routing decision: a direct exchange sends a message to the queue whose binding key matches the routing key, a topic exchange matches patterns, a fanout exchange ignores keys and copies to every bound queue. The queue is the buffer that holds messages until a consumer takes them, and it survives a consumer crash — which matters, because consumers crash. The binding is the rule connecting the two. None of this is hidden in raw AMQP. All of it is hidden behind @app.task.

Here is a consumer doing the two things that matter most: taking one message at a time, and acknowledging it only after the work is genuinely done.

import asyncio
import aio_pika

async def main() -> None:
    connection = await aio_pika.connect_robust("amqp://guest:guest@localhost/")
    channel = await connection.channel()

    # Take at most one unacknowledged message at a time.
    await channel.set_qos(prefetch_count=1)

    queue = await channel.declare_queue("work", durable=True)

    async with queue.iterator() as messages:
        async for message in messages:
            try:
                await handle(message.body)
                # Work succeeded: tell the broker it can forget this message.
                await message.ack()
            except Exception:
                # Work failed: requeue=False sends it to the dead-letter path
                # rather than spinning on the same poison message forever.
                await message.reject(requeue=False)

async def handle(body: bytes) -> None:
    ...  # your actual work

asyncio.run(main())

That is a complete, reliable worker in about a dozen meaningful lines. The reliability lives in two places: set_qos(prefetch_count=1) and the explicit ack/reject. Both are settings Celery makes on your behalf, with defaults you never chose and probably could not recite.

Acks, nacks, and the redelivery you signed up for

The acknowledgement is the load-bearing concept in the entire system, and it is the one Celery hides most thoroughly. When a consumer receives a message with manual acknowledgement on, the broker does not delete that message. It marks it as delivered-but-unacknowledged and keeps it, waiting. If the consumer acks, the broker drops it. If the consumer dies — process killed, network partition, machine rebooted — without acking, the broker notices the channel has gone and requeues the message for someone else to handle.

This is the mechanism behind at-least-once delivery, and it has a sharp edge. The broker cannot tell the difference between "the consumer died before doing the work" and "the consumer did the work and died before acking". From RabbitMQ's point of view those are identical: an unacknowledged message on a dead channel. So it redelivers. If your handler charged a card, sent an email, or incremented a counter before the ack, that side effect happens again.

Which is the whole game. At-least-once does not mean "RabbitMQ occasionally messes up". It means the contract is at least once, by design, and the obligation to make redelivery safe is yours. You handle that with idempotency — a dedupe key, an upsert keyed on the message's own identifier, a check-then-act guarded by the database. The broker gives you a guarantee with a known failure mode; it does not give you exactly-once, because exactly-once over an unreliable network is, in the general case, a myth. Anyone selling it to you has quietly redefined either "exactly" or "once".

Raw AMQP's other mode, autoack — where the broker considers a message delivered the instant it leaves the queue — throws all of this away. Crash mid-message and the message is simply gone. Celery does not use raw autoack, but don't mistake that for safety: in its default configuration (task_acks_late disabled) the worker acknowledges the message just before it starts executing the task. Crash mid-task and the broker has already been told the message is handled, so it is never redelivered — practically the same loss mode as autoack, arrived at by a different route. To get redelivery when a worker dies you have to turn on task_acks_late (and usually task_reject_on_worker_lost with it). The problem is not that Celery picks the wrong default; it is that it never makes you confront the choice, so you arrive at production having never decided what your system should do when a worker dies holding a message.

Prefetch, and why the default ruins throughput

Prefetch — the QoS setting — controls how many unacknowledged messages the broker will push to a single consumer at once. It is the least glamorous knob in RabbitMQ and the one most likely to be quietly destroying your throughput or your fairness, depending on which way it is wrong.

Set it too low, as in the snippet above with prefetch_count=1, and every consumer processes exactly one message before the broker sends the next. That is maximally fair — a slow message on one worker never blocks others — and for long, uneven tasks it is exactly right. But there is a round trip to the broker between every message, and if your messages are small and fast, that round trip dominates and throughput collapses. Set it too high, or leave it unbounded, and one greedy consumer hoovers a thousand messages into its local buffer, starving its peers and turning a clean restart into a thundering re-delivery of everything it was holding.

The right value depends on your message duration and processing variance, and you tune it by measuring rather than guessing — somewhere in the tens for short tasks, one or two for long ones. The point is not the number. The point is that there is a number, it has real consequences, and Celery sets it to a default you have never examined.

Dead-letter exchanges and the retry topology

The reject(requeue=False) in the consumer has to send the message somewhere, and that somewhere is a dead-letter exchange. You configure it when you declare the queue: a rejected or expired message, instead of vanishing, is republished to a nominated exchange, which routes it to a dead-letter queue. From there you can inspect poison messages by hand, or — more usefully — build a delayed retry by giving the dead-letter queue a message TTL and dead-lettering it back to the original. Message fails, lands in the retry queue, waits thirty seconds, returns to the work queue: a backoff loop assembled entirely from broker primitives, with no scheduler and no extra process.

You can see this whole arrangement at a glance:

rejectexpiresproducerwork exchangework queueconsumerdead-letter exchangeretry queue, TTL 30sparked: poisonmessagesrejectexpiresproducerwork exchangework queueconsumerdead-letter exchangeretry queue, TTL 30sparked: poisonmessages

It is a few lines of declaration, and once you have drawn it you understand your retry behaviour completely, because you built it out of parts whose behaviour you can name. Compare that to Celery's retry() with its max_retries and countdown, which works fine until you need to know where the message physically is during the backoff, and the answer turns out to involve the result backend, a separate piece of infrastructure you may not have realised you were depending on.

What you give up, and what you buy

Here is the trade stated plainly.

Celery hides this Doing it yourself buys you this
The broker — exchanges, queues, bindings A mental model of where every message physically is
The ack lifecycle and redelivery A clear answer to "what happens when a worker dies mid-message"
Prefetch / QoS tuning Throughput and fairness you set on purpose, not by default
Retry mechanics and the result backend A retry topology built from broker primitives, no hidden dependency
The exactly-once illusion Honest at-least-once, and idempotent handlers that make it safe

I am not going to pretend this is free. Celery gives you task chaining, a result backend, a beat scheduler, and a community of people who have hit your problem before. If you genuinely need workflow orchestration — fan-out, fan-in, chords, the whole DAG — then you are buying something real, and reimplementing it on bare AMQP is not a good use of your week. Use the tool that fits.

But most teams reaching for Celery are not orchestrating DAGs. They are running background work — send this email, resize this image, reconcile this account — and for that the framework is mostly hiding the very concepts they most need when the pager goes off at three in the morning. The cost of dropping it is a few dozen lines of connection setup and a willingness to learn five nouns. The return is that you can answer, with certainty, the question that ends most incidents: did the task run twice, and what did we do about it?

Own your delivery semantics. They are not someone else's to define.