Skip to content

Find which tasks wedge a Celery worker by reading kombu's Redis unacked hash

A Celery 5.6 worker (--pool=gevent --concurrency=10, task_acks_late=True, worker_prefetch_multiplier=1, Redis broker via kombu 5.6 over rediss://) stops consuming a few minutes to an hour after every restart. ECS reports the task RUNNING, the process uses almost no CPU, the queue list (LLEN <queue>) keeps growing, and nothing gets logged. celery inspect active isn't reachable because the worker runs --without-heartbeat --without-mingle --without-gossip, and there's no shell on the container to run py-spy. Redis looks healthy (no rejected connections, no evictions, memory low). How can you tell from the broker side which tasks are holding the worker's pool slots, without touching the queue?

1 solution
ranked by outcome — not votes
Accepted

With the Redis transport, kombu keeps every message that has been delivered to a consumer but not yet acked in two broker keys: a hash unacked (delivery_tag -> JSON [payload, exchange, routing_key]) and a sorted set unacked_index (delivery_tag scored by delivery time as a unix timestamp). With acks_late=True and prefetch 1, a message stays there until its task finishes. So when the entry count equals --concurrency and the entries never age out, the pool is exhausted by tasks that never return. That is your wedge, and the task names are the suspects.

Read-only inspection (redis-py 6.x):

import json, datetime, redis
r = redis.Redis.from_url(BROKER_URL)
for tag, ts in r.zrange('unacked_index', 0, -1, withscores=True):
    payload, exchange, rk = json.loads(r.hget('unacked', tag))
    h = payload['headers']
    print(datetime.datetime.fromtimestamp(ts, datetime.UTC).isoformat(), rk, h['task'], h.get('expires'))

In our case it showed exactly 10 entries (= concurrency), 8 of them the same 2-minute beat task delivered at each tick for about 12 minutes after the restart, and none of them ever acked. Every one went through a helper that runs asyncio.run() on a gevent.threadpool.ThreadPool thread. Also useful: CLIENT LIST shows how many consumers are sitting in cmd=brpop, and idle= on each one; a BRPOP connection idle for thousands of seconds is a dead consumer.

Caveats: keys live in the broker DB (db 0 by default) and may be prefixed if you set global_keyprefix. Don't HDEL/ZREM these to 'fix' things: kombu restores unacked messages to the queue after visibility_timeout. Also, the gevent pool doesn't enforce hard time_limit, so a hung task holds its slot until the process restarts. Put an explicit asyncio.wait_for/gevent.Timeout inside the task.