skip to content

questions

4

A service processes CPU-heavy tasks in a worker pool on an 8-core machine. A colleague proposes raising the pool from 8 workers to 200 to make it faster. Why does throughput usually not improve, and how can it get worse?

level: juniorimportance: must knowfreq 60%

answer

  1. concurrency != parallelism; cores cap parallelism
  2. context switch + cold cache + cold TLB
  3. throughput flattens at cores, then declines
  4. fixed rate + more in flight = proportionally more latency
  5. only blocking work wants many threads

basics

~20 s

Only 8 tasks can actually compute at once, so throughput is already capped by the cores. Extra workers add context switches, cache pollution, memory for stacks and more lock contention. Throughput flattens then dips, and per-task latency rises because more work is in flight.

solid answer

~50 s

Concurrency is not parallelism. With CPU-bound work the machine executes at most one task per core at a time, so once the pool reaches the core count the CPU is saturated and more workers cannot raise the completion rate. What extra workers do add is overhead: context switches costing microseconds each, eviction of each task's working set from cache when it is descheduled, longer scheduler run queues, memory per stack, and more contenders on any shared lock. Throughput therefore flattens near the core count and then declines. Latency behaves worse than throughput: with the completion rate fixed, doubling the number of tasks in flight roughly doubles the time each spends in the system, so a change meant to 'do more at once' makes every request slower. The queue in front of the pool is not waste - it holds the same pending work far more cheaply than a parked thread. More threads help only when workers spend most of their time blocked rather than computing.

go deeper

for a junior

Say that only as many tasks as cores can actually run at once and that extra threads cost context switching and memory.

for a middle

Add cache and TLB pollution and lock contention, and describe the throughput curve peaking near the core count.

for a senior

Bring in the latency consequence of more work in flight at a fixed completion rate, and redirect the conversation to identifying the bottleneck resource.

for a principal

Frame it as capacity management: pool size is an admission-control decision, and once a resource is saturated the levers are cheaper work, load shedding, or more capacity - not more threads.

## Concurrency versus parallelism Concurrency is how many tasks are in progress; parallelism is how many execute simultaneously. Parallelism is bounded by hardware - roughly the number of hardware threads. Raising the pool size raises concurrency, not parallelism. For work that is genuinely computing, the extra concurrency buys nothing because the bottleneck resource is already fully occupied. ## Where the extra time goes **Context switching.** Each involuntary switch costs a direct kernel cost of a few microseconds plus a much larger indirect cost: the incoming task finds cold caches and a cold TLB and runs slowly until its working set is reloaded. With 200 runnable workers each gets a short slice and is descheduled mid-computation, so a large share of every slice is spent re-warming. **Cache capacity.** Eight tasks may fit their hot data in the shared last-level cache; two hundred will not. The workload shifts from cache-resident to memory-bound, which can cut per-task speed by a large factor - this is how throughput ends up strictly worse rather than merely flat. **Memory.** Each thread reserves stack space and the runtime keeps per-thread bookkeeping. Hundreds are affordable; the habit that follows ('just raise it again') is not. **Contention.** If tasks share any lock, adding contenders lengthens the queue behind it and can trigger convoying, where every worker serializes behind the same lock while still paying scheduling costs. ## The shape of the curve Throughput rises nearly linearly to the core count, flattens, and then declines as overhead grows. The peak sits at or slightly above the core count for CPU-bound work - slightly above because even compute-heavy tasks take page faults and occasional I/O, so a little oversubscription keeps cores busy. ``` throughput | ____ | / \____ | / \____ | / +--------------------------- workers ^ cores ``` ## Latency is the part people forget In a stable system the average number of tasks present equals the completion rate times the average time each spends there. If the completion rate is pinned by the cores, increasing the number of tasks in flight increases their average residence time proportionally. Twenty-five times more workers on a saturated 8-core box means roughly twenty-five times the latency per task with no extra work done. Deadlines and timeouts then start firing, retries pile on, and the system degrades faster than the raw overhead alone would suggest. ## When more threads genuinely help If a task spends most of its time blocked - waiting on a network call, a disk, a database - it holds a worker without holding a core. Then the pool must be much larger than the core count so some worker is always runnable. The correct size is a function of the ratio of waiting to computing, which is exactly why mixing the two kinds of work in one pool makes the number impossible to choose. ## The right instinct Identify the bottleneck resource first. If CPU utilization is already near saturation, more workers cannot help and the levers are cheaper work, load shedding, or more machines. If CPU is idle while requests queue, the bottleneck is elsewhere - a blocking dependency, a lock, a connection pool - and enlarging the thread pool usually just moves the queue to a more expensive place, or overwhelms the downstream system that was the real constraint.

  • When would a pool much larger than the core count be the right answer?
    When workers spend most of their time blocked rather than computing - synchronous network calls, disk reads, database round trips. A blocked worker holds no core, so many are needed to keep the cores busy. The size then follows the ratio of waiting time to service time, and must also respect the downstream system's capacity, since a large pool can simply overload the dependency you were waiting on.
  • Your CPU-bound pool sits at the core count and requests still queue. What do you do?
    Accept that the machine is saturated and stop tuning the pool. The options are to reduce the work per request (better algorithms, caching, cheaper serialization), shed or throttle load so latency stays bounded, or add machines and spread the work. Enlarging the pool would only lengthen residence time and hurt the tail while completing no more work per second.

A kitchen with 8 stoves. Hiring 200 cooks does not cook more food; they crowd the aisles, bump into each other, and every dish takes longer to reach the pass.

saying these in an interview costs you the question

  • Treating thread count as a throughput dial independent of the hardware
  • Ignoring latency because throughput 'stayed the same'
  • Believing a queued task is more expensive than a parked thread
  • Assuming a context switch costs only the kernel's direct cost, not the cache and TLB damage
  • Sizing without first knowing whether the work is CPU-bound or blocking

context

open as a page

Derive a sizing rule for a worker pool from first principles: how many workers does a purely CPU-bound workload need, and how does that number change when each task spends most of its time waiting on something external?

level: middleimportance: must knowfreq 55%

basics

~20 s

A worker is runnable only the fraction S/(S+W) of its life, where S is compute time and W is wait time. To keep C cores busy at target utilization U you need N = C x U x (1 + W/S) workers. CPU-bound means W is zero, so N is about the core count.

open as a page

Explain Little's law and how you would use it, together with observed utilization and queue depth, both to choose the number of workers for a queue-driven service and to explain why response time degrades sharply as the pool approaches saturation.

level: seniorimportance: should knowfreq 42%

basics

~20 s

Little's law: in a stable system the average number of items present equals arrival rate times average time in the system, L = lambda x W. Busy workers needed is arrival rate times service time. Response time grows like 1/(1 - utilization), so the last few percent cost enormous latency.

open as a page

A service runs three kinds of work on one shared worker pool: fast in-memory lookups, slow calls to a third-party HTTP API, and periodic report generation. Make the case for splitting them into separate pools, and explain how you would allocate capacity across the pools.

level: principalimportance: should knowfreq 45%

basics

~20 s

One pool cannot have one correct size: the three classes have wildly different wait/service ratios, and slow tasks occupy workers so fast ones queue behind them. Split into pools sized per class, budget CPU across them, cap each class's concurrency, and measure utilization and queue wait per pool.

open as a page