Celery acks_late and Worker Crash Safety
By default Celery acknowledges a task message before running it, which means a worker that crashes mid-task loses that task for good. This guide switches to late acknowledgement and configures the settings that make it work correctly, as part of Celery Architecture & Configuration in Backend Frameworks & Worker Scaling.
Problem Statement
A data platform runs Celery on Kubernetes with the RabbitMQ broker and default settings. Nodes are occasionally reclaimed, and workers are OOM-killed a few times a week on large inputs. A reconciliation audit found that about 0.1% of "import file" tasks never completed and never failed — they simply vanished, leaving files unimported with no error logged. Every one coincided with a worker being killed. You want tasks that are interrupted by a crash to run again, prefetched tasks not to be lost or stuck, clarity on which tasks must now be idempotent, and a test that proves it.
Prerequisites
- Celery 5.3+ with RabbitMQ or Redis as broker.
- Tasks that are idempotent or can be made so (every task that is redelivered may run twice).
- Knowledge of your longest task duration (for the Redis visibility timeout).
- A staging environment where you can kill worker processes.
Step 1 — See When Celery Acknowledges
With the default early acknowledgement, the worker acks the message as soon as it receives it and hands it to a pool process. From that moment the broker considers the task done. If the pool process dies, the task is gone.
default (acks_early): receive -> ACK -> run task -> store result
^ broker forgets the message here
task_acks_late=True: receive -> run task -> ACK
^ broker forgets only after success/failure
Early ack gives at-most-once execution; late ack gives at-least-once. Which one is right is a per-task decision, but for most tasks — imports, emails, syncs — losing work silently is worse than occasionally doing it twice. Delivery semantics in general are covered in Exactly-Once vs At-Least-Once Delivery.
Step 2 — Enable Late Ack and Reject on Worker Loss
Two settings work together. task_acks_late moves the ack after execution. task_reject_on_worker_lost controls what happens when the pool process executing a task dies while the worker's main process survives (the OOM-killer often kills just the child): with it set, the message is rejected and requeued; without it, Celery marks the task failed with WorkerLostError and acks it.
# celeryconfig.py
task_acks_late = True
task_reject_on_worker_lost = True # requeue when a pool child dies mid-task
worker_prefetch_multiplier = 1 # see Step 3
task_acks_on_failure_or_timeout = True # still ack tasks that raise: they go through retry logic
Per-task overrides are possible for tasks that must never run twice and are cheap to lose — @app.task(acks_late=False) — but make that the documented exception rather than the default.
Step 3 — Reduce Prefetch So Crashes Release Little
With late ack, each worker holds concurrency × prefetch_multiplier unacked messages. When a worker dies, all of them are redelivered — including ones it had not started. With the default multiplier of 4 and concurrency 8, a crash releases 32 messages, and a slow task can hold 31 others hostage in the prefetch buffer while it runs.
worker_prefetch_multiplier = 1 # at most one reserved task per pool process
# For very short tasks where throughput matters more, 2-4 is reasonable;
# for long tasks, 1 prevents head-of-line blocking behind a slow task.
There is a second, quieter cost of a high multiplier with late ack. Prefetched messages are reserved to one worker even while it is busy, so a worker stuck on a slow task holds several quick tasks that other, idle workers could have run. With a multiplier of 1, each pool process reserves only the task it is about to start, and work flows to whichever process is free.
Prefetch sizing is covered in depth in tuning prefetch and consumer concurrency.
Step 4 — Set the Redis Visibility Timeout (Redis Broker Only)
On RabbitMQ, an unacked message is redelivered when the consumer's connection closes. On Redis, Celery emulates this with a visibility timeout: a message not acked within visibility_timeout seconds is redelivered to another worker — even if the original worker is still running it. The default is one hour.
broker_transport_options = {
"visibility_timeout": 7200, # > your longest task, including countdown/eta delays
}
result_backend_transport_options = {"visibility_timeout": 7200}
Set it above the longest task duration and the longest countdown/eta you use, because Redis-backed Celery also holds scheduled tasks as unacked messages. Too short, and long tasks run twice concurrently; too long, and crashed tasks wait that long to be retried. The trade-off is the same one described in the visibility timeout deep dive.
Step 5 — Make Redelivered Tasks Safe
Late acknowledgement means every task can run more than once. Audit tasks for effects that must not repeat and make them idempotent: upserts instead of inserts, a processed-marker checked in the same transaction as the effect, idempotency keys on external calls.
@app.task(bind=True, acks_late=True)
def import_file(self, file_id: int):
with db.transaction():
f = File.objects.select_for_update().get(pk=file_id)
if f.status == "imported":
return "already imported" # redelivery after a crash post-commit
rows = parse(storage.read(f.key))
Row.objects.bulk_create(rows, ignore_conflicts=True) # unique (file_id, line_no)
f.status = "imported"
f.save(update_fields=["status"])
The unique constraint on (file_id, line_no) makes a partially-completed import safe to rerun; the status check makes a completed one a no-op. Techniques are covered in preventing duplicate job execution with idempotency.
Step 6 — Test with SIGKILL
Prove the behaviour by killing a pool process mid-task, in the same way the OOM killer would.
# Start a worker with one process; enqueue a slow task; kill the child mid-run
celery -A app worker -c 1 --loglevel=info &
python -c 'from tasks import import_file; import_file.delay(42)'
sleep 5
pkill -9 -f "celery.*ForkPoolWorker" --newest # kill the pool child, not the main process
# Expect in the worker log: task requeued (not "WorkerLostError ... failed")
# Expect the task to run again and complete; file 42 imported exactly once
Automate this as described in chaos testing worker crashes and redelivery, and run it for every task type with irreversible effects.
Verification
After enabling late ack in production, the reconciliation audit should find zero "vanished" tasks. Also watch redelivery volume so that a misconfigured visibility timeout does not go unnoticed:
# Tasks started more than once (from a counter incremented when request.delivery_info.redelivered)
sum(rate(celery_task_redelivered_total[1h])) by (task)
A steady low rate matches crashes and deploys; a rate tied to long tasks points at a visibility timeout that is too short.
Gotchas & Edge Cases
Tasks that kill workers every time. With reject_on_worker_lost, a task that always OOM-kills its worker is requeued forever. Track redeliveries per task id and send repeat offenders to a dead-letter queue; quorum queues in RabbitMQ can enforce a delivery limit.
Graceful shutdown still matters. Late ack protects against crashes; SIGTERM handling decides what happens on deploys. See handling SIGTERM in Celery workers.
eta tasks on Redis. Long countdowns occupy the visibility timeout budget; very long delays belong in a scheduler, not countdown.
Pool choice. With the gevent or eventlet pools, there are no child processes to lose — a crash takes the whole worker, and every unacked message is redelivered. task_reject_on_worker_lost matters mainly for the prefork pool; late ack itself matters for all pools. Pool trade-offs are in choosing Celery prefork vs gevent pools.
Chords and late ack. A chord member redelivered after it already counted toward the chord can make the body fire twice. Make chord bodies idempotent as well.
Result backend writes. A task can complete its work and crash before storing its result; the redelivered run must produce the same result.
FAQ
Why isn't acks_late the default? Because it requires idempotent tasks, and Celery cannot know whether yours are. Early ack is the conservative default for duplicate-sensitive code, at the cost of silent loss.
Does acks_late help if the broker restarts? Unacked messages survive a broker restart only if queues and messages are durable (persistent delivery mode on RabbitMQ, AOF on Redis).
How do I roll this out safely?
Audit tasks first: list every task, mark whether a second execution is harmless, and fix or opt out the ones that are not. Enable task_acks_late for workers of one low-risk queue, watch redelivery counts and duplicate-effect checks for a week, then extend queue by queue. Opt-outs (acks_late=False) should come with a comment explaining why losing the task is acceptable.
Should failed tasks be acked?
Yes — task_acks_on_failure_or_timeout=True (the default) acks tasks that raise, so failures go through retry and result logic instead of being redelivered endlessly.
Related
- Celery Architecture & Configuration — Celery's components and settings.
- Celery Task Time Limits — bounding tasks that would otherwise run forever.
- Fixing Celery Worker Memory Leaks — removing the OOM kills in the first place.
- RabbitMQ Consumer Timeout for Unacked Messages — the broker-side limit on unacked tasks.