Guide · RabbitMQ

Celery and RabbitMQ for LLM Task Queues

Many Python AI teams already run Celery for background work, and RabbitMQ is one of the brokers Celery's documentation lists as stable, with monitoring and remote control supported. LLM jobs stress that pairing in ways ordinary web tasks do not: one call can take minutes, providers throttle under load, and every retry costs tokens. This guide covers the settings that matter, in the order they cause incidents, from the Celery 5.6 and RabbitMQ 4.3 documentation checked on 10 October 2026.

Tyler Eastridge

By Tyler Eastridge, Head of Operations

LinkedIn · Updated

8 min read13 sections
On this page
The setup in one paragraph

The setup in one paragraph

Set task_acks_late = True and worker_prefetch_multiplier = 1, so each worker process holds one job and a crashed worker's job returns to the queue. Keep the hard task time limit below the RabbitMQ consumer timeout, 30 minutes by default, and raise that timeout per queue. Use quorum queues with publisher confirms. Retry rate-limit errors with exponential backoff, and make every task idempotent, because some jobs will run twice. Use Redis or a database as the result backend, route GPU and API-bound work to separate queues, and alert on unacknowledged messages and consumer counts.

Why LLM work belongs on a queue

A queue takes slow, expensive model calls off the request path, caps how many run at once (the number of worker processes becomes your concurrency limit against the provider), and gives every job a retry policy. Celery supplies task definitions, retries, workflows such as groups and chords, and result storage. RabbitMQ supplies delivery: a job leaves the queue only when the worker acknowledges it. Celery's broker overview notes that RabbitMQ as the broker with Redis as the result backend is a very common pairing, and this guide assumes that layout. Most of the trouble comes from defaults chosen for short tasks.

Acknowledge late and prefetch one

By default Celery acknowledges a message just before running the task, so a deploy or an out-of-memory kill halfway through a ten-minute call drops the job. Celery's optimization guide gives the setting for long-running tasks: task_acks_late = True with worker_prefetch_multiplier = 1. The worker then reserves only as many tasks as it has processes and acknowledges each after it returns. At the default multiplier of 4, the first worker to start can hoard a backlog of minute-long jobs while others sit idle.

Two details catch people out. A late-acknowledged task that raises an exception or hits its time limit is still acknowledged by default (task_acks_on_failure_or_timeout), so failures are handled by retries, not redelivery. And a worker process killed by a signal is acknowledged too, unless task_reject_on_worker_lost is on, which Celery warns can cause message loops.

Keep task time limits under the consumer timeout

RabbitMQ gives up on a consumer that holds a delivery unacknowledged for longer than the delivery acknowledgement timeout, 30 minutes by default, and returns its messages to the queue. With late acknowledgement the clock runs for the whole LLM call, so a 40-minute batch job is requeued and started elsewhere while the first copy is still running. Celery's documentation warns of the same PRECONDITION_FAILED error for long countdowns.

From RabbitMQ 3.12 a queue can carry its own timeout through the consumer-timeout policy key or the x-consumer-timeout queue argument, in milliseconds. RabbitMQ 4.3 applies acknowledgement timeouts only to quorum queues and lets a consumer set its own with an x-consumer-timeout consumer argument; a classic queue on 4.3 has no timeout, so a hung worker keeps its jobs. Use a policy, because redeclaring an existing queue with different arguments fails with PRECONDITION_FAILED, while a policy can change at any time:

# One policy per queue applies, so set the timeout and dead-lettering together
rabbitmqctl set_policy llm-queues "^llm\." \
  '{"consumer-timeout": 900000, "dead-letter-exchange": "llm.dlx"}' \
  --apply-to quorum_queues

The order to keep: soft limit, then hard limit, then consumer timeout, each comfortably above the last.

A configuration to start from

Steps 2 and 3 plus quorum queues, routing, retries and an idempotency check, in one file. LLM_URL stands for whatever model API you call.

# celery_app.py: Celery 5.6 on RabbitMQ 4.3 quorum queues
import hashlib, os

import httpx, redis
from celery import Celery
from kombu import Queue

app = Celery("llm", broker=os.environ["BROKER_URL"], backend=os.environ["RESULT_URL"])
app.conf.update(
    task_acks_late=True,
    worker_prefetch_multiplier=1,
    broker_transport_options={"confirm_publish": True},  # required for quorum queues
    task_queues=[
        Queue("llm.api", queue_arguments={"x-queue-type": "quorum"}),
        Queue("llm.gpu", queue_arguments={"x-queue-type": "quorum"}),
    ],
    task_default_queue="llm.api",
    task_routes={"embeddings.*": {"queue": "llm.gpu"}},
    task_soft_time_limit=600,
    task_time_limit=660,  # below the queue's consumer timeout
)
cache = redis.Redis.from_url(os.environ["CACHE_URL"], decode_responses=True)


class RetryableLLMError(Exception):
    pass  # a 429 or 5xx from the model endpoint


@app.task(
    autoretry_for=(RetryableLLMError, httpx.TransportError),
    retry_backoff=2, retry_backoff_max=300, retry_jitter=True, max_retries=8,
)
def summarize(doc_id: str, text: str, model: str) -> str:
    key = f"summary:{doc_id}:{model}:{hashlib.sha256(text.encode()).hexdigest()[:16]}"
    if (done := cache.get(key)) is not None:
        return done  # a redelivered task costs no tokens
    resp = httpx.post(os.environ["LLM_URL"], json={"model": model, "input": text}, timeout=300)
    if resp.status_code == 429 or resp.status_code >= 500:
        raise RetryableLLMError(resp.status_code)
    resp.raise_for_status()  # other 4xx errors fail without retrying
    summary = resp.json()["output"]
    cache.set(key, summary, ex=7 * 86400)
    return summary

Use quorum queues, and know what changes

Celery supports quorum queues from 5.5: declare queues with x-queue-type: quorum or set task_default_queue_type = "quorum", and turn on publisher confirms, which Celery's documentation says they require. A quorum queue's contents are agreed by a majority of its replicas, so losing one node of three loses no queued jobs.

Three behaviors change. Quorum queues do not support global QoS, so Celery turns it off when it detects them; prefetch becomes static and worker autoscaling stops working. Countdown and ETA tasks would then block a worker slot, so Celery switches to native delayed delivery and schedules them inside RabbitMQ, which covers every retry with backoff. And from RabbitMQ 4.0 a quorum queue drops or dead-letters a message after 20 failed deliveries by default, so give each queue a dead-letter target.

Retry rate limits with backoff

autoretry_for lists the exceptions to retry. retry_backoff=True waits 1, 2, 4, 8 seconds and so on, and a number sets the factor. retry_backoff_max caps the wait at 600 seconds by default, retry_jitter (on by default) randomizes each delay so throttled tasks do not retry in lockstep, and max_retries defaults to 3, low for a busy provider. Retry a 429 or a 5xx, never a prompt that exceeds the context window, which fails the same way and costs tokens every time. Celery's rate_limit is enforced per worker instance, not across the cluster; for a hard provider limit, give that model its own queue and size its workers.

Make every task idempotent

Delivery is at least once. A consumer timeout, a connection lost while a late-acknowledged task runs, and task_reject_on_worker_lost all lead to a second run, sometimes in parallel with the first. Celery plans to cancel such tasks on connection loss by default in 6.0 (worker_cancel_long_running_tasks_on_connection_loss), and turning it on now is the safer choice for long jobs. Build a key from document ID, content hash and model version, check for a stored result before calling the model, and write with an upsert. Embedding pipelines get most of this free when each vector has a deterministic ID.

Batch embeddings with chords, and pick the result backend

A group runs tasks in parallel, chunks splits a long argument list into fixed-size tasks, and a chord runs a callback once a whole group has finished, for example to mark an index build complete. Chords need a result backend, and their tasks must not ignore results. Celery's rpc:// backend returns results as AMQP messages: a result can be read once, only by the client that sent the task, and is transient unless result_persistent is set, and chords are not supported on it. Use Redis or SQL for chords, status dashboards and agent runs other services read. Send document IDs, not megabytes of text, in messages.

Set time limits, and pick the pool with them in mind

task_soft_time_limit raises SoftTimeLimitExceeded so the task can clean up; task_time_limit kills and replaces the worker process. Give the HTTP client a timeout below the soft limit. Prefork, the default pool, supports both limits. Celery's documentation notes that other pools silently disable soft time limits, and gevent does not enforce the hard limit on a blocking task, so API-bound tasks moved to gevent or eventlet rely on client timeouts.

Route by model, GPU and run time

Celery's optimization guide recommends separate workers for long and short tasks, because settings that suit one hurt the other. A common split: local inference and embedding models on a queue only GPU hosts consume, with concurrency matched to GPU memory; hosted-API calls on a high-concurrency queue; post-processing on the default queue. task_routes maps task names, with wildcards, to queues, and celery worker -Q picks the queues a worker consumes. Split by provider too, so one throttled endpoint does not stall the other's jobs.

Monitor in-flight work, not just depth

Flower is the monitor Celery's documentation recommends for task states and workers. RabbitMQ's Prometheus plugin exposes broker metrics on port 15692. For LLM queues, watch unacknowledged messages per queue (jobs in flight), the consumer count (a drop to zero while workers run usually means a consumer timeout), the redelivery rate (duplicate spend) and arrivals on the dead-letter queue.

Running LangGraph or LangChain agents from Celery

LangGraph and LangChain do not document a Celery integration; LangChain's hosted Agent Server runs its own task queue on PostgreSQL and Redis. Teams self-hosting a graph often use this pattern, which is common practice, not an official integration: the Celery task takes a thread ID, the graph uses a persistent checkpointer such as PostgreSQL (an in-memory one is lost on restart), and LangGraph's documented way to resume after an error is to invoke with None and the same thread ID:

@app.task(bind=True, acks_late=True, soft_time_limit=600, time_limit=660)
def run_agent(self, thread_id: str, user_input: dict) -> dict:
    config = {"configurable": {"thread_id": thread_id}}
    snapshot = graph.get_state(config)   # graph compiled with a Postgres checkpointer
    payload = None if snapshot.next else user_input   # pending nodes: resume, do not restart
    return graph.invoke(payload, config, durability="sync")

Avoid the "exit" durability mode here: it saves state only when a run ends and cannot recover from a crash midway. A failed node runs again on resume, so tool calls with side effects still need idempotency keys.

The pitfalls that show up most

  • Prefetch left at 4 with minute-long jobs.
  • A hard time limit above the consumer timeout, so long jobs run twice.
  • Classic queues on RabbitMQ 4.3, where a hung worker holds its jobs indefinitely.
  • Countdown tasks on classic queues, held unacknowledged in worker memory until due.
  • autoretry_for=(Exception,), paying for retries that can never succeed.
  • No dead-letter target on quorum queues, so a message that fails 20 deliveries is dropped.

These failures sit between Celery and the broker. AceMQ provides RabbitMQ support for Celery workloads, and RabbitMQ vs Kafka vs NATS for AI agents covers the choice of broker.

Frequently asked questions

What prefetch setting should Celery use for LLM tasks?

worker_prefetch_multiplier = 1 with task_acks_late = True. Each process reserves one task and acknowledges it when finished, so long jobs spread evenly and a crashed worker's job returns to the queue.

Why do my Celery tasks run twice on RabbitMQ?

Usually a late-acknowledged task outran the consumer timeout, 30 minutes by default, or the worker lost its connection mid-task, and RabbitMQ requeued it. Keep the hard time limit below the queue's consumer timeout and make tasks idempotent.

Does Celery support RabbitMQ quorum queues?

Yes, from Celery 5.5, with confirm_publish enabled. Celery disables global QoS when it detects quorum queues and uses native delayed delivery for countdown and ETA tasks.

Can I run LangGraph agents with Celery?

Yes, as a pattern rather than an official integration: invoke the graph inside a task with a thread ID and a persistent checkpointer, and resume a redelivered run by invoking with None and the same thread ID.

RabbitMQ services

Where this gets done

The work behind this page, run by the same engineers who wrote it.

More resources

Other RabbitMQ guides, comparisons and research

From the blog

Recent RabbitMQ articles

Next step

Need this done on your cluster?

AceMQ's senior RabbitMQ engineers support 130+ enterprise clients in 26+ countries under a 15-minute emergency SLA, with direct escalation to the RabbitMQ core team.