On this page
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_queuesThe 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 summaryUse 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.
Related
Where this gets done
The work behind this page, run by the same engineers who wrote it.
- 24/7 RabbitMQ support15-minute emergency SLA, versions back to 3.8.x
- Managed RabbitMQ servicesWe run the brokers, on your infrastructure or hosted
- RabbitMQ consultingArchitecture, migration and remediation from senior engineers
- RabbitMQ health checkEngineer-led assessment with a prioritised fix list
- Extended LTS support for RabbitMQ 3.xCVE backports for versions the community no longer patches
- RabbitMQ commercial licensingTanzu RabbitMQ licences from an authorized Broadcom partner
- RabbitMQ troubleshootingLive incidents and recurring faults
- RabbitMQ upgrades3.x to 4.x, planned and executed in your window
- RabbitMQ migrationsFrom IBM MQ, Kafka, cloud brokers or older RabbitMQ
- RabbitMQ implementation and architectureCluster design, DR and go-live
- RabbitMQ corporate trainingAdmin and developer courses taught by working engineers
Other RabbitMQ guides, comparisons and research
- GuideThe RabbitMQ Performance Tuning GuideRead the guide
- GuideThe RabbitMQ Streams GuideRead the guide
- GuideThe RabbitMQ Reliability Guide: Ten Failure Patterns and Their FixesRead the guide
- GuideThe RabbitMQ Disaster Recovery GuideRead the guide
- GuideThe RabbitMQ Clustering and Sizing GuideRead the guide
- GuideThe RabbitMQ on Kubernetes GuideRead the guide
- GuideThe RabbitMQ Migration GuideRead the guide
- GuideThe RabbitMQ Security and Hardening GuideRead the guide
- GuideAgent-to-Agent Messaging over AMQP and RabbitMQRead the guide
- ComparisonRabbitMQ vs Amazon SQS ComparedSee the comparison
- ComparisonRabbitMQ vs Redis ComparedSee the comparison
- ComparisonManaged RabbitMQ Options ComparedSee the comparison
- ComparisonMessage Broker Support Options ComparedSee the comparison
- ResearchWhat Breaks in Production RabbitMQ: 145 Support Tickets, 2023 to 2026Read the research
- ResearchRabbitMQ in Production 2026: What 22 Assessed Estates Actually RunRead the research
- ResearchThe RabbitMQ CVE Register, 2026 EditionRead the research
Recent RabbitMQ articles
- RabbitMQ Cluster Operator vs Helm Chart on KubernetesOct 2026
- RabbitMQ Alternatives: 9 Options and When to StayOct 2026
- RabbitMQ End of Life and End of Support Dates (3.6 to 4.3)Oct 2026
- VMware Licensing Cost in 2026Sep 2026
- The Tanzu Software in Your VCF You Are Not UsingSep 2026
- Production RabbitMQ Architecture for Enterprise TeamsSep 2026
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.