Right-Sizing Worker Concurrency per CPU
Before you add more worker replicas, make sure each replica is doing as much work as its CPU allows. Concurrency per worker (Celery's --concurrency, BullMQ's concurrency, Sidekiq's thread count) is the setting that decides this, and it is usually picked by guesswork. This guide derives it from measurements, as part of Horizontal Worker Scaling in Backend Frameworks & Worker Scaling.
Problem Statement
A team runs Celery workers on Kubernetes pods with 2 vCPU each and the default prefork concurrency of 2. One queue sends transactional email through an HTTP API; another resizes images. The email queue pods sit at 8% CPU while its backlog grows, so the autoscaler adds pods — 40 of them — and the Redis connection count climbs toward its limit. Someone raises concurrency to 64 on all queues, and the image queue immediately slows down: jobs take four times longer, the pods are CPU-throttled, and memory triples. You want one method that gives the right concurrency for each queue, with numbers you can defend, instead of a single value applied everywhere.
Prerequisites
- Per-job timing for each queue: total duration, and ideally CPU time (from
resource.getrusage,process.cpuUsage(), or profiling). - Container CPU requests and limits you control.
- A staging environment or a quiet production window where you can load-test one worker.
- Worker metrics: throughput (jobs/s), CPU utilisation, memory, and throttling (
container_cpu_cfs_throttled_periods_total).
Step 1 — Measure CPU Time vs Wait Time per Job
The key number is the ratio between the time a job spends running on a CPU and the time it spends waiting (network calls, database queries, disk). Record both for each job:
import resource, time
from celery.signals import task_prerun, task_postrun
_start = {}
@task_prerun.connect
def _before(task_id=None, **_):
ru = resource.getrusage(resource.RUSAGE_THREAD)
_start[task_id] = (time.perf_counter(), ru.ru_utime + ru.ru_stime)
@task_postrun.connect
def _after(task_id=None, task=None, **_):
wall0, cpu0 = _start.pop(task_id, (None, None))
if wall0 is None:
return
ru = resource.getrusage(resource.RUSAGE_THREAD)
wall = time.perf_counter() - wall0
cpu = (ru.ru_utime + ru.ru_stime) - cpu0
JOB_WALL.labels(task.name).observe(wall)
JOB_CPU.labels(task.name).observe(cpu)
In Node.js, process.cpuUsage() is per process, so measure it with concurrency set to 1 during a test run. Collect a few thousand jobs and take the median and p90 of both values per job type. For the email job you might see 400 ms wall time with 12 ms CPU; for image resizing, 900 ms wall with 850 ms CPU.
Step 2 — Compute a Starting Concurrency
If a job uses the CPU for a fraction u of its wall time, one core can in theory keep 1/u such jobs busy. Aim below 100% so there is room for spikes and the runtime itself:
concurrency_per_core ≈ target_utilisation / cpu_fraction
cpu_fraction = cpu_time / wall_time
email: 0.7 / (12 / 400) = 0.7 / 0.03 ≈ 23 per core
image: 0.7 / (850 / 900) = 0.7 / 0.94 ≈ 0.75 per core → 1 per core
For a 2-vCPU pod, that suggests roughly 46 concurrent email jobs and 2 concurrent image jobs. The same result follows from Little's law: throughput per pod equals concurrency divided by wall time, and CPU caps throughput at cores × target / cpu_time. Setting the two equal gives the formula above.
Treat this as a starting point. The formula assumes wait time is really idle; a job waiting on a lock in your own database is not free, because more concurrent jobs make that wait longer.
Step 3 — Pick the Concurrency Model That Fits
The execution model limits which numbers are practical:
| Runtime | Model | Practical per-core range |
|---|---|---|
| Celery prefork | one process per slot | 1–4 (memory-bound) |
| Celery gevent / eventlet | green threads | 20–200 for I/O jobs |
| RQ | one process per worker | 1–2 |
| BullMQ | async on one event loop | 10–100 for I/O jobs; 1 for CPU |
| Sidekiq | threads under the GIL | 5–25 |
| Go workers | goroutines | limited by downstream, not CPU |
A Celery queue that needs 46 concurrent I/O jobs should not run 46 prefork processes — each process carries its own copy of the application (often 150–300 MB). Run that queue under gevent, or use a threads pool, and keep prefork for the CPU-bound queue. In BullMQ, CPU-bound work blocks the event loop, so give it concurrency: 1 per process and use sandboxed processors or more processes to use additional cores.
Step 4 — Split Queues by Profile
A single concurrency setting cannot serve jobs that differ by 30× in CPU fraction. Route them to separate queues and run a worker deployment per profile:
# I/O-bound: many green threads per pod, small CPU request
celery -A app worker -Q email,webhooks -P gevent -c 48
# CPU-bound: one process per core, CPU request = limit
celery -A app worker -Q images,reports -P prefork -c 2
The split also makes autoscaling honest: the I/O pool scales on backlog, and the CPU pool's CPU utilisation actually reflects its load. See dedicated worker pools for routing patterns.
Step 5 — Align Container CPU Limits with Concurrency
Kubernetes CPU limits are enforced by CFS quota in 100 ms periods. A pod with a 2-CPU limit running 4 CPU-bound processes gets throttled every period, and job durations stretch unpredictably. For CPU-bound pools:
- Set concurrency equal to the CPU request, and set the limit equal to the request (or drop the limit and rely on requests for scheduling).
- Watch
container_cpu_cfs_throttled_periods_total / container_cpu_cfs_periods_total; above 5–10% means the pod wants more CPU than it may use. - Make the runtime aware of the limit: Node.js and Python report the host's core count, not the container's, so
os.cpus().lengthoros.cpu_count()can suggest 64 on a pod limited to 2. Derive concurrency from an environment variable set alongside the limit.
For I/O pools, the CPU request can be small, but memory grows with concurrency: each in-flight job holds its payload, response buffers, and connections. Measure memory per in-flight job and set the memory limit for the chosen concurrency plus headroom.
Step 6 — Load-Test One Worker and Find the Knee
The formula gives a starting value; the load test finds the real one. Fill the queue with representative jobs, run a single worker pod, and step concurrency up while recording throughput, p95 duration, and CPU:
for c in 8 16 32 48 64 96; do
kubectl set env deploy/worker-io CONCURRENCY=$c
kubectl rollout status deploy/worker-io
sleep 300 # let throughput settle
./record-metrics.sh "c=$c"
done
Choose the value just before the knee, where throughput stops rising and latency starts climbing. If the knee comes much earlier than the formula predicted, the bottleneck is not your CPU — look at the downstream service's rate limit, the database connection pool, or lock contention. That is also the signal that more replicas will not help either; see rate limiting & throttling jobs.
Verification
- Throughput per pod at the chosen setting is close to the plateau of the load test.
- CPU utilisation for CPU-bound pools sits at 60–80% under load, with throttling below 5%.
- p95 job duration under full load is within ~20% of the single-job duration.
- The replica count needed for peak load dropped (the email pool above went from 40 pods to 3), and broker connection counts dropped with it.
- Memory per pod stays below 80% of the limit at full concurrency.
Gotchas & Edge Cases
Hidden CPU in "I/O" jobs. JSON parsing of large responses, TLS handshakes, and template rendering add CPU that the job author does not think of as work. Measure; do not assume.
Connection pools must match. Forty-eight concurrent jobs sharing a database pool of 10 will queue on the pool, and the measured wait time will include that queue. Size pools to concurrency, or cap concurrency to the pool.
Downstream limits are shared. Concurrency per pod multiplied by replicas is the load the downstream sees. Raising per-pod concurrency and keeping the autoscaler maximum unchanged can multiply the peak load on a third-party API.
Prefetch interacts with concurrency. In Celery, worker_prefetch_multiplier × concurrency messages are reserved per worker. With 48 green threads, a multiplier of 4 reserves 192 messages that other pods cannot take. Lower it to 1 for long jobs.
Hyperthreads are not cores. A vCPU on most clouds is one hyperthread. CPU-bound work gains perhaps 20–30% from the second hyperthread, not 100%.
FAQ
Should concurrency equal the number of cores? Only for CPU-bound jobs. For I/O-bound jobs, concurrency equal to cores leaves the CPU mostly idle and forces you to scale out replicas instead; the right value is often 10–50× the core count.
Why did raising concurrency make every job slower? The jobs were CPU-bound, or they share a bottleneck (database, lock, rate-limited API). More concurrency only queues them for that resource inside the worker instead of in the broker.
Is it better to run more small pods or fewer large ones? For I/O-bound pools, fewer pods with higher concurrency use fewer connections and less memory overhead. For CPU-bound pools, small pods (1–2 CPU) scale in finer steps. Either way, fix concurrency per pod first, then scale replicas.
How often should I re-measure? Whenever the job code or a dependency changes meaningfully, and at least quarterly. A new library version or a slower downstream API changes the CPU fraction and moves the knee.
Related
- Horizontal Worker Scaling — scaling principles.
- Configuring BullMQ Concurrency Limits for High Throughput — BullMQ-specific settings.
- Scaling Workers with KEDA on Queue Length — scaling replicas once each pod is sized.
- Forecasting Backlog Drain Time — using per-pod throughput in capacity plans.