skip to content

A batch computation was parallelised across all available cores and ended up slower than the single-threaded version, with correct results. Walk through the causes you would investigate and how you would distinguish between them.

level: seniorimportance: must knowfreq 46%

answer

  1. correct but slower = performance bug, not a race
  2. flat curve → overhead; declining → contention; early plateau → saturated resource
  3. false sharing: distinct variables, same cache line
  4. per-worker locals, merge once
  5. parallelism multiplies compute only

basics

~20 s

Look for: per-unit scheduling overhead exceeding the work, contention on a shared lock or counter, false sharing of cache lines between workers, saturation of a shared resource (memory bandwidth, one disk, one connection pool), too little total work to amortise startup, and oversubscription causing context switches. Distinguish by checking whether speedup falls as workers are added — that points to contention or bandwidth, not overhead.

solid answer

~1 min

Group the causes by their signature. **Overhead-dominated** — units too small, so submit/join costs exceed the work; also per-item allocation or boxing driving garbage collection. Signature: total time roughly constant or mildly worse regardless of worker count; huge allocation or task-queue counts. Fix: coarser units, a sequential cutoff. **Contention-dominated** — a shared lock, a shared counter, or an accumulator that every worker touches; also **false sharing**, where workers write distinct variables sitting on the same cache line. Signature: throughput *degrades* as workers increase; high system time or cache-coherence traffic. Fix: per-worker local state merged once; pad or block-align hot output. **Bandwidth/resource-bound** — the work is memory-streaming, or all workers hit one disk, one connection pool, or a rate-limited service. Signature: speedup plateaus after a few workers no matter what; low IPC, high memory stall or I/O wait. Fix: reduce data movement, batch, or cap concurrency to what the resource supports. **Structural** — input too small for startup cost, a serial split or merge step, non-random-access input that is expensive to partition, or oversubscription against other work on the box. Method: measure sequential baseline, sweep worker count 1→N and read the shape of the curve, then confirm with allocation, lock and cache counters before changing anything.

code

text · 11 lines
text
workers:      1     2     4     8     16

A)  time:   100    98    99   101   104     flat/slightly worse
            => per-unit overhead dominates (units too small, allocation)

B)  time:   100    70    85   120   190     improves then degrades
            => contention: shared lock/counter, or false sharing

C)  time:   100    60    45    44    45     early plateau
            => saturated shared resource: memory bandwidth, one disk,
               a small connection pool, a rate limit

go deeper

for a junior

Name the basic reasons: the pieces are too small to be worth scheduling, there is not enough total work, or the threads are fighting over one shared thing.

for a middle

Separate overhead from contention, know what false sharing is, and describe the per-worker-local-accumulator fix.

for a senior

Drive it from measurement: a warmed baseline, a worker-count sweep whose shape names the cause, then counters for allocation, locks, cache invalidation and memory stalls before changing anything.

for a principal

Frame it at system level — parallelism multiplies compute only; in an already-concurrent service intra-request parallelism often trades throughput for latency, and CPU quotas, NUMA and shared pools decide the real ceiling.

## Why this happens at all Parallel execution adds costs that the sequential version never pays: task creation, scheduling, synchronisation, cache-coherence traffic, and contention on any shared hardware resource. When those costs exceed the arithmetic saved by running on multiple cores, the parallel version loses. Correct results and worse performance is the normal signature — this is a performance bug, not a correctness one. ## Cause 1: per-unit overhead exceeds per-unit work If you create a task per element and the element takes hundreds of nanoseconds while the task envelope costs a microsecond, most of the machine is doing bookkeeping. Related: per-item object allocation, boxing of primitives, and closure creation, which turn a tight loop into allocation pressure and garbage-collection work that did not exist before. *Signature*: worker-count sweep shows roughly flat or slightly worsening time; profilers show scheduler internals and allocation near the top. *Fix*: bigger units, a sequential cutoff, chunking so each unit is worth tens of microseconds at minimum. ## Cause 2: contention on shared mutable state Every worker increments a counter, appends to a shared list, updates a shared cache, or writes to a shared logger. The mutual exclusion serialises the exact code every worker runs, so you get Amdahl-shaped behaviour plus coordination overhead. Even lock-free atomics contend — a hot atomic counter under 32 cores can be slower than a mutex, because the cache line holding it bounces between all of them. *Signature*: performance gets **worse** with more workers — the tell-tale sign of contention rather than overhead. High context-switch counts or spin time. *Fix*: per-worker local accumulators merged once at the end; sharded counters; read-mostly structures; removing the shared logger or cache from the hot path. ## Cause 3: false sharing Different workers write *different* variables that happen to occupy the same cache line (typically 64 bytes). No logical sharing exists, but the hardware coherence protocol invalidates the whole line on every write, so it ping-pongs between cores. Common when workers write into adjacent slots of a results array, or when per-worker state objects are allocated next to each other. *Signature*: dramatic slowdown that disappears when you pad the per-worker data or give each worker a contiguous block instead of interleaved indices; hardware counters show cache-line invalidations. *Fix*: pad per-worker state to a cache line, aggregate into a local variable and write once at the end, or switch from interleaved to block partitioning. ## Cause 4: a saturated shared resource The work is not CPU-bound at all. Streaming through a large array is limited by **memory bandwidth**: a handful of cores can saturate the memory bus, so cores 5 through 32 add nothing and cost coherence traffic. Similarly: all workers reading from one disk (and on spinning media, turning sequential reads into seeks), all calling one downstream service, all drawing from a connection pool of size four, all hitting a rate limiter. *Signature*: speedup rises to some small number of workers and then flatly plateaus or declines; low instructions-per-cycle with high memory-stall or I/O-wait time. *Fix*: reduce bytes moved (better layout, compression, operating on smaller types), improve locality so the data stays in cache, batch I/O, or simply cap concurrency at the level the resource supports — extra workers past that point are pure cost. ## Cause 5: not enough work Starting workers, warming pools, partitioning and merging all cost time. For a millisecond of work, that overhead is the whole runtime. A cost-based check before parallelising ('if estimated work < threshold, run sequentially') is standard in library implementations for exactly this reason. ## Cause 6: serial fractions you forgot The split step may be expensive (partitioning a non-random-access structure requires traversal), the merge may be serial (sorting or deduplicating combined results on one thread), and setup/teardown is serial by definition. If half the runtime is a serial merge, the ceiling is low no matter how many cores work on the parallel half. ## Cause 7: oversubscription and environment Running P workers when the process is already handling concurrent requests, or when it is confined to fewer CPUs than the machine reports (a container CPU quota), means threads fight for time slices and each pays context-switch and cache-eviction costs. A common production trap: a parallel routine used inside a request handler multiplies threads by concurrent requests, degrading total throughput even when each single request looks faster in isolation. Under NUMA, workers on one socket touching memory allocated on another add remote-access latency. ## A diagnostic order that works 1. **Establish a real sequential baseline**, warmed up, on the same data. 2. **Sweep the worker count** 1, 2, 4, 8, … and plot. The shape names the cause: flat from the start → overhead; rises then declines → contention or false sharing; rises then plateaus early → a saturated shared resource. 3. **Check the obvious structural things**: is the input big enough? Is there a shared accumulator in the loop body? Is the merge serial? 4. **Bring in counters**: allocation rate, context switches, lock wait time, cache misses and invalidations, memory bandwidth, I/O wait. 5. **Change one thing and re-measure.** Coarsen granularity; replace shared state with per-worker locals; pad hot per-worker data; cap concurrency at the resource limit. 6. **Accept the honest outcome**: some workloads are memory-bound or too small, and the right answer is to keep them sequential and optimise the data layout or the algorithm instead. The general lesson worth stating explicitly: parallelism multiplies *compute*, and only compute. If the bottleneck is memory traffic, a device, a downstream service, or a lock, adding workers moves the queue rather than shortening it.

  • You sweep the worker count and see the time improve from 1 to 2 workers, then get steadily worse from 4 to 16. What does that shape tell you?
    Degradation with added workers is the signature of contention rather than fixed overhead: a shared lock, a hot atomic, or false sharing on a cache line, where the cost grows with the number of participants. Fixed per-task overhead would show a flat curve instead, and a saturated resource would plateau rather than decline. Next step is to look for shared mutable state in the loop body and at cache-invalidation counters.
  • How would you tell a memory-bandwidth limit apart from lock contention?
    Bandwidth limits plateau — throughput rises to a few cores and then stops improving without getting much worse — and show low instructions-per-cycle with high memory-stall time and a memory bus near its measured ceiling. Contention actively degrades with more workers and shows lock wait time, context switches, or cache-line invalidations on a specific address. A quick check is to shrink the working set so it fits in cache: if the scaling suddenly appears, it was bandwidth.
  • Why can a parallel routine that measurably speeds up one request reduce a service's overall throughput?
    Because the service is already parallel across concurrent requests, so the cores are not idle. Parallelising inside a request multiplies threads by in-flight requests, causing oversubscription, context switching and cache eviction, plus contention on any shared pool. Single-request latency improves in isolation, but total throughput and tail latency get worse under load. Container CPU quotas make this worse, since the process may see far fewer CPUs than the host reports.

saying these in an interview costs you the question

  • Assuming a correct-but-slower parallel version must contain a race; it is a performance problem, not a correctness one.
  • Believing more threads always help, without checking whether the bottleneck is memory bandwidth, a device, or a lock.
  • Overlooking false sharing because the code has no logically shared variables.
  • Comparing against a cold or unwarmed sequential baseline, or measuring on data far smaller than production.
  • Adding parallelism inside a request handler in an already-concurrent service and judging it by single-request latency alone.

context