Architecture¶
Worker model¶
The celery-asyncio worker uses a multi-threaded architecture:
Main thread
|
+-- Consumer event loop (asyncio)
| - Broker connection (Redis/AMQP)
| - Message dispatch
| - Timer/scheduler
|
+-- Loop worker threads (N)
| - Each runs its own asyncio event loop
| - Async tasks dispatched here
| - Semaphore limits concurrency per thread
|
+-- Sync worker threads (M)
- ThreadPoolExecutor
- Sync tasks dispatched here
Main thread¶
Owns the broker connection and runs the consumer event loop. Messages are received here and dispatched to the appropriate worker thread based on whether the task is async or sync.
Loop workers¶
Each loop worker runs its own asyncio event loop. Async tasks (async def) are scheduled directly on these loops. A semaphore controls how many tasks can run concurrently per loop worker.
The number of loop workers and per-worker concurrency are configurable:
# 2 loop workers, 500 concurrent tasks each = 1000 total
celery -A app worker --loop-workers=2 --loop-concurrency=500
Sync workers¶
Synchronous tasks run in a ThreadPoolExecutor. The number of sync worker threads is configurable:
Python 3.14t free-threading¶
With Python 3.14t (free-threaded build), all threads run with true parallelism since the GIL is disabled. This means sync tasks in the thread pool actually execute in parallel, not just concurrently.
Consumer loop¶
The consumer loop in the main thread:
- Drains the timer heap (fires scheduled entries like ETA/countdown tasks)
- Blocks on
drain_events()until a message arrives or timeout - Batch-drains remaining messages non-blocking to fill the concurrency pipeline
- Checks restart conditions (max tasks, max memory, stuck threads)
Bootsteps¶
The worker startup/shutdown sequence uses async bootsteps. Each step's start() and stop() methods are coroutines:
Shutdown and restart¶
- Signals:
SIGTERMand the firstSIGINTstart a warm shutdown.SIGQUITor a secondSIGINTstarts a cold one, and so doesSIGTERMwithREMAP_SIGTERM=SIGQUIT.SIGHUPrestarts a worker detached from its terminal, after a warm shutdown. - Warm shutdown: the worker stops consuming and drops the prefetched tasks that have not started, leaving their messages for the broker to redeliver. Running async tasks get
worker_soft_shutdown_timeoutseconds to finish (0 by default). The ones still running after that are cancelled as a cold shutdown cancels them. Sync tasks are left to finish. - Cold shutdown: the worker drops the prefetched tasks as a warm shutdown does, and cancels the running tasks immediately, sync ones as far as their thread can be stopped. An
acks_latetask is left unacknowledged for the broker to redeliver. Any other task was acknowledged when it started, so it is stored asRETRYand does not run again. A task whose body has returned is left to report its result. - Sync tasks: a cold shutdown and
revoke(terminate=True)stop a sync task by raising an exception in its thread, only while the task body runs. Python code stops immediately. A call blocked in C code, such astime.sleepor a socket read, sees the exception when it returns, and the worker waits for that before it closes the broker connection, as it waits for the sync tasks a warm shutdown leaves to finish. - Draining:
worker_max_tasks_per_child,worker_max_memory_per_childor a stuck thread makes the worker stop consuming and finish the tasks it already holds, prefetched ones included. When every sync thread is stuck, a task that has not started cannot start, so the worker finishes only the started ones and leaves the others for the broker to redeliver. It then runs its exit handlers, which save the--statedbfile, and restarts itself withos.execv. - Hard time limit: an async task is cancelled, and a sync task is stopped as above. Either is reported as failed. A sync task's thread still running 2 seconds later counts as stuck. Neither the restart that follows nor a shutdown waits for that thread, which the exec or the exit ends.