HotShard
Routing under load

Thread Pools & Work Queues

How many threads should this pool have, and what happens when they're all busy?

A server that starts a fresh thread for every task falls over twice. First on the cost of creating threads faster than it retires them. Then on a traffic spike that creates so many threads the machine runs out of memory. A thread pool fixes both by fixing a number: a set of worker threads created once, pulling tasks off a shared queue. That leaves two questions, and they are the whole topic. How many workers? And what happens when they are all busy?

~6 min read

Start here: why not a thread per task?#

TL;DRthe 30-second version
  • A thread pool is a fixed set of reusable worker threads pulling tasks from a shared queue. You neither create a thread per task (expensive) nor let the thread count grow without bound (a memory and context-switch blow-up).
  • Sizing is the core question. CPU-bound work wants about one thread per core. I/O-bound work wants many more, because a thread blocked on I/O holds no core.
  • The queue in front of the pool must be bounded. An unbounded queue hides overload until the process is killed for running out of memory.

The simplest way to run work in the background is to start a thread for each piece of it. A request arrives, you spawn a thread, it does the work, it exits. This works right up until load arrives. A thread is not free. Each one reserves a stack (commonly around a megabyte) and takes a kernel scheduling slot. Under a steady stream of short tasks you spend a large share of your time making and tearing down threads instead of doing the work.

The second failure is worse. A thread-per-task design has no ceiling. If tasks arrive faster than they finish, the thread count climbs without limit. Ten thousand threads is gigabytes of stacks and ten thousand things for the scheduler to juggle, so the cores spend their time context-switching instead of computing. The machine thrashes and stops.

The mechanism: a queue and N workers#

A thread pool has two parts: a set of worker threads, and a task queue they share. You submit a task, a function to run, by putting it on the queue. Each worker runs the same tiny loop forever: take the next task off the queue, run it to completion, come back for another. When the queue is empty the workers sleep on it at zero CPU cost, and one is woken when a task arrives.

submittersmany callers
enqueue task
task queuebounded, FIFO
take next
N workersreused, capped
run to completion
doneworker loops back
Submit to the queue; workers drain it
A worker holds a thread, not a coreN workers does not mean N cores busy. A worker running CPU work occupies a core. A worker blocked on a disk read or a network reply holds its thread but no core, because the OS has parked it and given the core to someone else. That one distinction is why the right pool size depends on whether your tasks compute or wait.

The sizing question: how many workers?#

  • CPU-bound tasks (hashing, compression, parsing) keep a core busy the whole time. More workers than cores does not help. The extra threads take turns on the same cores and add context-switch overhead. Target roughly the core count, often core count plus one to cover the brief stalls on a cache miss or page fault.
  • I/O-bound tasks (a database call, an HTTP request, reading a file) spend most of their time blocked. Here you want many more workers than cores. While one worker waits on I/O, its core is free for another worker that has computing to do.
The pool-sizing formula (Goetz, Java Concurrency in Practice)Nthreads = Ncpu × Ucpu × (1 + W/C). Ncpu is the number of cores, Ucpu is your target CPU utilization (0 to 1), and W/C is the ratio of time a task spends waiting to time it spends computing. For pure CPU work W is about 0, so you get about one thread per core. For a task that waits 90 ms on I/O and computes for 10 ms, W/C = 9, so on 8 cores at full utilization you want 8 × (1 + 9) = 80 threads.
PredictA service calls a downstream API that takes about 100 ms to reply, of which only ~5 ms is CPU on your side. You run 4 cores. A junior engineer sets the pool to 4 threads 'because 4 cores' and throughput is terrible. What size does the workload actually want, and why?

Hint: How much of each task is actually on a core, and what is the thread doing the rest of the time?

Far more than 4, on the order of 80. Each task spends ~95 ms blocked and only ~5 ms computing, so W/C ≈ 19. With 4 threads, almost all of them are parked waiting on the API and your cores sit idle. The formula gives 4 × 1 × (1 + 19) = 80.

When the pool is full: bound the queue, then choose#

It is tempting to use an unbounded queue so no task is ever refused. That is a trap. An unbounded queue adds no workers, so the drain rate is unchanged. The backlog grows, the wait grows with it, memory grows, and the process is eventually OOM-killed. That is a hard crash instead of a clean refusal.

So the queue must be bounded. Once it is, you have to answer the real question: when the queue is full and every worker is busy, what happens to the next task? This is the saturation policy. Java's ThreadPoolExecutor names the standard choices:

PolicyWhat it does when the queue is fullWhen to use it
Abort (reject)Throw an error immediately; the submitter must handle itThe default. Fail fast and let the caller retry, back off, or shed. Honest load shedding.
Caller-runsRun the task on the submitting thread itselfNatural backpressure: the submitter is busy running the task, so it cannot submit more. Slows the producer without dropping work.
BlockMake the submitter wait until a queue slot freesProducers you control. Pushes the slowdown upstream, but risks stalling the caller.
Discard / discard-oldestSilently drop the new task (or the oldest queued one)Rarely. A task that vanishes without a trace is the hardest overload to debug.

If this comes up in an interview#

The one-linerA fixed set of reused workers draining a bounded queue. Size it from the work: about the core count for CPU-bound, far more for I/O-bound. When the queue fills, reject or push back, never grow without limit.
Is a bigger pool always better under load?

No. Past the size the workload needs, more threads add context-switch overhead and memory pressure without adding throughput. For CPU-bound work that point is around the core count. An image-thumbnailing service on 8 cores with a pool of 500 threads gets slower under a spike, not faster.

One slow dependency made my whole service unresponsive. Why?

Almost certainly a shared pool. One downstream hangs, its calls hold every worker, and requests that never touch that dependency can't get a thread either. The fix is isolation: give each risky dependency its own bounded pool, so its failure exhausts only its own workers. That is the bulkhead pattern, and it is how Netflix's Hystrix contained failures. Pair it with a circuit breaker so you stop calling a dead dependency at all.

When would you use an event loop instead of a pool?

When concurrency is I/O-heavy. A pool pays a parked thread, about 1 MB of stack, per waiting task. An event loop holds an idle wait for a few bytes of kernel state, but it runs on one core and one long CPU task freezes every other connection. Strong systems use both. Node's own runtime keeps a four-thread pool for file and DNS calls so the JavaScript thread never blocks.

Common mistakes & gotchas
  • Pool tasks that wait on other tasks in the same pool. If every worker is busy waiting on a task still sitting in the queue behind them, nothing ever completes. This is thread-pool starvation. Give dependent tasks separate pools, or use a work-stealing pool built for sub-tasks (Java's ForkJoinPool).
  • Sizing from a guess. The number comes from the work: CPU-bound wants about the core count, I/O-bound wants Goetz's formula or Little's Law.
Go deeperUnder the hood: the ThreadPoolExecutor core/max/queue quirk

Java's ThreadPoolExecutor has four knobs: corePoolSize (threads it keeps alive even when idle), maximumPoolSize (the hard ceiling), the workQueue, and the rejection handler. A new task creates a thread up to corePoolSize. Beyond that the task is queued. Only when the queue is full does the pool create threads up to maximumPoolSize. Only when the queue is full and the maximum is reached does the rejection handler fire.

So if you pair a large maximumPoolSize with an unbounded queue, the pool never grows past corePoolSize, because the queue never fills. The max you configured is dead code. If you want the pool to scale up under load, the queue must be bounded. Core, max, queue, and rejection are one coupled decision, not four.

References & further reading
References

Feedback on this topic →