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.

CPU time vs wait time per job An email job takes 400 milliseconds, of which only 12 milliseconds use the CPU; the rest is waiting for the HTTP API. An image resize job takes 900 milliseconds, of which 850 are CPU. The first can run dozens per core; the second about one per core. Where a job's wall time goes email send waiting on HTTP API (388 ms) 400 ms wall, 12 ms CPU → 3% CPU image resize on CPU (850 ms) 900 ms wall, 850 ms CPU → 94% CPU CPU waiting (I/O)

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
Splitting work by CPU profile Jobs from the application are routed by type: email and webhook jobs go to an I/O pool running 48 green threads on half a CPU, while image and report jobs go to a CPU pool running 2 processes on 2 dedicated cores. Each pool scales on its own backlog. One deployment per job profile application routes by job type email, webhooks 3% CPU per job images, reports 94% CPU per job I/O pool: gevent -c 48 0.5 CPU request CPU pool: prefork -c 2 2 CPU request = limit Each pool autoscales on its own backlog.

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().length or os.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
Finding the concurrency knee As concurrency per pod increases from 8 to 96, throughput rises nearly linearly up to 48 and then flattens. p95 job duration is flat until 48 and then climbs steeply, because jobs start queuing for CPU or a downstream resource. The right setting is just before the knee. Throughput and p95 duration vs concurrency 8 16 32 48 64 96 concurrency per pod knee throughput p95 duration

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