How would you size and split a Celery worker fleet running CPU-heavy report rendering and thousands of webhook calls: pools, concurrency, prefetch and autoscale?
answer
- two workloads, two worker types
- cores or RAM, whichever runs out
- downstream limits bound greenlets
- autoscale resizes one pool only
basics
~20 sSplit Celery into two worker deployments on separate queues: prefork reports with concurrency bounded by cores and RAM, prefetch 1 and child recycling; gevent webhooks with hundreds of greenlets and a larger prefetch. Scale mostly by adding workers.
solid answer
~40 sI would run two worker types consuming separate queues. Report workers: `-P prefork`, `-c` = the smaller of the core count and RAM divided by a report's peak memory, `worker_prefetch_multiplier = 1` so long tasks never hoard short ones, and `worker_max_memory_per_child` as a guard rail. Webhook workers: `-P gevent -c 200` or so, bounded by what receivers and our outbound connection pool tolerate, with a larger multiplier because tasks are short. `--autoscale=max,min` only resizes the pool inside one worker between two numbers, and its prefetch window is sized from the maximum; I would reserve it for bursty load on a fixed host and scale containers on queue depth otherwise. The trade-off is a second deployment to operate in exchange for isolation.
go deeper
Recall that CPU-heavy and I/O-bound tasks want different Celery pools, so they usually run on different workers.
Explain how cores, RAM and prefetch each limit report workers, and why greenlet counts depend on downstream limits instead.
Size each worker type from measured peaks and latency, set prefetch per workload, and explain autoscale's prefetch and host limits.
Own the call on when splitting pays for its operational cost, which metric shows the line has moved, and how the fleet scales on queue depth.
## Start from the two workloads Sizing a **Celery** fleet starts from what each task spends its time on, because the **pool** (`-P`) and **concurrency** (`-c`) are properties of a worker, not of a task. | | Report rendering | Webhook calls | |---|---|---| | Bottleneck | CPU and memory | waiting on remote HTTP | | Run time | minutes | hundreds of milliseconds | | Pool | `prefork` | `gevent` (or `eventlet`) | | Concurrency bound | cores and RAM | receivers and connection limits | | Prefetch | multiplier 1 | multiplier 4 or more | Putting both on one worker forces one answer to every row. On a gevent worker a report freezes every webhook; on a prefork worker each waiting HTTP call occupies a whole process, and long reports hoard prefetched webhooks. ## Sizing the report workers 1. Measure a report's **peak resident memory** and run time on realistic data. 2. Set `-c` to the smaller of the **core count** and **usable RAM ÷ peak per report**. On an 8-core, 16 GB host with 6 GB peaks, that is 2, not 8. 3. Set `worker_prefetch_multiplier = 1` so a worker never reserves more than one extra report per child. 4. Add `worker_max_memory_per_child` (kilobytes) so a child that has grown past a chosen size is replaced after its task, shedding the high-water mark. 5. If step 2 leaves cores idle, the task is memory-bound: chunk it or buy memory-heavy hosts. ## Sizing the webhook workers Cores are not the limit for greenlets; almost all their time is spent waiting. The real bounds are: - how many **simultaneous requests** the receivers tolerate before they throttle or fail; - the size of your **HTTP client's connection pool** — greenlets beyond it just queue for a connection; - **timeouts** on every call, so a hung receiver cannot occupy slots indefinitely. Start around `-c 100`–`300`, watch latency and error rates, and move from there. Because these tasks are short, a larger prefetch multiplier buys throughput without starving anything. ## Autoscale versus more workers `--autoscale=10,3` keeps at least 3 slots and grows towards 10 when the worker has received more tasks than it has pool slots, shrinking again some time after load drops. The workers guide lists it for prefork and gevent. Its limits shape the decision: - it resizes **one pool on one host**; the host's CPU and RAM do not grow with it; - the prefetch window is computed from the **maximum** concurrency, so a worker at 3 children may already hold `10 × multiplier` messages; - for report workers, growing past the RAM-based `-c` is exactly the failure you sized against. Autoscale fits a bursty queue on a fixed machine. On a container platform, scaling the **number of workers** on queue depth usually gives cleaner isolation and predictable per-worker sizing. ## A starting configuration For the fleet in question, a defensible first version looks like this, to be tuned against measurements: - **reports**: `-P prefork -c 2` on 16 GB hosts, `worker_prefetch_multiplier = 1`, `worker_max_memory_per_child` below the 6 GB peak (for example `3000000`, about 3 GB) so a child that just rendered a large report is replaced — a restart is cheap next to a minutes-long task, scaled out by adding workers when report queue wait grows; - **webhooks**: `-P gevent -c 200`, default multiplier 4, strict HTTP timeouts, scaled out when webhook queue wait or call latency grows; - **both**: one queue per workload, so neither worker ever reserves the other's tasks. ## When one fleet is fine, and what splitting costs Splitting is not free: a second deployment, a second set of dashboards, idle headroom on each side, and routing configuration that must stay in sync. One mixed worker is defensible when: - both workloads are small and latency on webhooks does not matter; - reports are short enough that prefetch starvation is invisible. Once reports take minutes or webhook latency is promised to a customer, the split pays for itself. The judgement a lead owns is where that line sits for this system, and which metric — queue wait time per queue — will tell the team when it has moved.
- When is Celery's --autoscale the right tool rather than running more worker containers?When load is bursty on a machine whose size is fixed, such as a single VM, and idle slots between bursts would waste memory. It resizes the pool within one worker between a max and a min. Where you can add containers on queue depth, that usually scales better and keeps each worker's size predictable.
- How would you choose the gevent concurrency for the Celery webhook workers?Start from the receivers' tolerance and the HTTP client's connection-pool size, not from cores. Set explicit timeouts, begin around a few hundred greenlets, and raise or lower while watching latency, error rates and time spent waiting for a connection.
saying these in an interview costs you the question
- One worker type with a large --concurrency serves both workloads equally well.
- Report worker concurrency should always equal the core count, whatever each task's memory.
- Celery's --autoscale adds new worker machines when the queue grows.
- Webhook worker concurrency should match the core count because each call is a task.