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.

Ack timing decides what a crash loses With early acknowledgement, the worker acks on receipt, starts the import, and is OOM-killed; the broker has already forgotten the message, so the import is lost silently. With late acknowledgement, the ack is sent only after the task finishes; when the worker dies, the message is still unacked and the broker redelivers it to another worker. Worker OOM-killed mid-task acks early ACK running import killed: task lost, no error acks late running import (unacked) killed: broker redelivers Late ack converts silent loss into a visible, idempotency-dependent retry.

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 decides the blast radius A worker with concurrency 8 and prefetch multiplier 4 reserves 32 unacknowledged messages. When it crashes, all 32 are redelivered, including 24 it had not started, and while it runs a slow task those reserved quick tasks wait. With a multiplier of 1 it reserves only 8, one per pool process, so a crash releases 8 and no idle capacity elsewhere is starved. Unacked messages held by one worker (concurrency 8) multiplier 4 8 running 24 reserved, waiting: redelivered on crash multiplier 1 8 running nothing else held back Lower prefetch costs a little throughput on tiny tasks and buys fairness and smaller crash impact.

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.

Redis visibility timeout cuts both ways With a two-hour visibility timeout, a task from a worker that crashed is redelivered two hours later. A healthy task that runs for three hours exceeds the timeout, so the broker redelivers it to a second worker at the two-hour mark while the first is still running, and the task executes twice concurrently. visibility_timeout = 2 h crashed task waits out the timeout redelivered at 2 h 3 h task worker A running worker B starts same task 2 h

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