skip to content

How do you choose task granularity — the amount of work per parallel unit — when each unit carries a fixed scheduling overhead? Explain the cost model and what goes wrong at each extreme.

level: seniorimportance: must knowfreq 48%

answer

  1. overhead ≪ per-unit work AND units ≫ workers
  2. unit work ≥ overhead / acceptable-fraction
  3. n ≈ 4–16 × P for balance; tail = one unit
  4. sequential cutoff; run one half inline
  5. empty interval = don't parallelise

basics

~30 s

Each unit costs a fixed overhead (submit, schedule, synchronise) on top of its useful work. Make units large enough that overhead is a small fraction of the work — typically at least tens of microseconds of work per unit — but small enough that you have several times more units than workers, so the tail after the last split is short. Too fine: overhead dominates. Too coarse: idle workers and a long tail.

solid answer

~1 min

Model it as `total = (work / P) + n · overhead + tail`, where n is the number of units. **Lower bound on unit size (efficiency).** Each unit pays a fixed cost: enqueue, hand-off, possible steal, result join, allocation of the unit itself — on modern hardware roughly hundreds of nanoseconds to a few microseconds. Keep that under ~1–10% of the unit's useful work, which puts the practical floor at tens of microseconds of computation per unit (often quoted as ~10,000+ basic operations). **Upper bound on unit size (balance).** Wall clock is the slowest worker, so the tail after the last unit is handed out is bounded by one unit's duration. With exactly P units, one slow unit wastes up to (P−1) worker-seconds. Aim for several times more units than workers — 4× to 16× P is the usual guidance — so the scheduler can even things out. So: `overhead ≪ per-unit work` and `units ≫ workers` simultaneously. If both cannot hold, the input is too small to parallelise, which is itself the answer. In practice you implement this as a **sequential cutoff**: split recursively while the range is large, and run the remainder sequentially below a threshold you measure rather than guess.

code

text · 14 lines
text
measured per-task overhead   c = 1 microsecond
acceptable overhead fraction f = 5%
workers                      P = 8
total work                   W = 200 ms

efficiency floor : work per unit >= c / f = 20 microseconds
                   => n <= W / 20us = 10,000 units
balance floor    : n >= 8 * P = 64 units

valid range: 64 .. 10,000 units  -> pick ~500 units (~400us each)

if instead W = 300 microseconds:
   balance wants n >= 64  -> ~4.7us per unit  -> overhead is ~21%
   the interval is empty  -> run it sequentially

go deeper

for a junior

Know that splitting work into very many tiny pieces costs more in scheduling than it saves, and that pieces should be reasonably chunky.

for a middle

State both constraints — overhead must be a small fraction of per-unit work, and there must be more units than workers — and explain the sequential cutoff in recursive splitting.

for a senior

Work the arithmetic from a measured overhead, discuss what happens at each extreme, describe how you would measure the cutoff, and recognise the case where the input is too small to parallelise at all.

for a principal

Argue for adaptive granularity (split-on-demand, guided tapering) over fixed constants, and tie the choice to the scheduler's redistribution ability and to the variance of unit costs in production data.

## What granularity means **Granularity** is the amount of useful work in one schedulable unit. Fine-grained means many small units; coarse-grained means few large ones. The choice is not stylistic: it determines whether parallelising helps at all, because every unit carries a fixed cost that does no useful work. ## What the overhead actually consists of Per unit, independent of its size, you pay some subset of: - allocating or boxing the unit (a task object, a closure, a future); - pushing it onto a queue and, later, popping it — usually an atomic operation, sometimes a contended one; - the scheduler's bookkeeping, and a possible hand-off to another core; - cache effects — a unit that runs on a different core than the one that produced its data pays cold-cache and cross-socket costs; - synchronisation on completion — a join, a counter decrement, a wake-up; - for a memory-bound unit, the cost of losing whatever the producing core had warm. Order of magnitude: sub-microsecond for a well-tuned in-process scheduler, several microseconds if a hand-off and a wake-up occur, and milliseconds if the "unit" crosses a process or network boundary. The absolute number matters less than the ratio. ## The cost model With total work W, P workers, n units and per-unit overhead c: ``` T_parallel ≈ W/P + n·c/P + T_tail + T_serial ``` The overhead term grows linearly with n; the imbalance term shrinks as n grows. That opposition is the entire tradeoff. **Efficiency constraint.** You want overhead to be a small fraction f of useful work: ``` c / (W/n) ≤ f → W/n ≥ c/f ``` With c = 1 µs and f = 1%, each unit needs ≥ 100 µs of work. With f = 10%, ≥ 10 µs. This is the origin of the common rule of thumb "a task should be worth at least tens of microseconds" or "at least ~10,000 operations". **Balance constraint.** The tail is bounded by the duration of the last unit still running, so: ``` T_tail ≈ W/n → n ≥ k·P with k ≈ 4…16 ``` Over-partitioning is what lets a scheduler (particularly a work-stealing one) redistribute when unit costs turn out unequal. With n = P exactly, one unit that is twice as expensive as the rest doubles your wall clock. Combine them: `k·P ≤ n ≤ W·f/c`. If the interval is empty — that is, `k·P·c/f > W` — the input is too small to parallelise profitably. That is a legitimate and often correct answer. ## What goes wrong at each extreme **Too fine.** Overhead dominates and the parallel version can be slower than the sequential one, sometimes by a large multiple. Symptoms: high CPU with low useful throughput, scheduler queues hot, allocation rate spiking, and speedup that *decreases* as you add workers because contention on the shared queue rises. The pathological case is a per-element task for elements that take nanoseconds — you have wrapped a 2 ns operation in a 1 µs envelope. **Too coarse.** Workers finish early and idle. With 8 workers and 8 units where one takes twice as long, utilisation drops to about 56%. Coarse units also react badly to any unpredictability — a unit that hits a slow path, a cache miss storm, or a preempted thread cannot be redistributed once it has started. ## The sequential cutoff The standard implementation of "right granularity" in recursive decomposition is a threshold: ``` solve(range): if size(range) <= CUTOFF: return solve_sequentially(range) (left, right) = split(range) a = spawn solve(left) b = solve(right) # run one half on the current worker return combine(await a, b) ``` Two details matter. First, running one half inline on the current worker halves the number of scheduled units for free. Second, the cutoff should be expressed in *work*, not item count, when per-item cost varies — 1,000 short strings and 1,000 megabyte blobs are not the same unit. How to pick CUTOFF: measure. Time the sequential version at several input sizes, find the size at which one unit takes roughly 50–100× your measured per-task overhead, and check that the resulting unit count is still several times the worker count. Then sweep the value and look at the curve — it is usually flat over a wide range, which is good news: you need to be in the right order of magnitude, not exact. ## Adaptive alternatives Rather than a fixed cutoff, some schedulers decide dynamically: split only if another worker is idle (or if the local queue is empty), so granularity adapts to available parallelism at runtime. Guided chunking, which hands out progressively smaller chunks as remaining work shrinks, is the data-parallel analogue: coarse early for efficiency, fine late for balance. ## What to say in an interview Name both constraints, give the ratio reasoning rather than a magic number, mention that the number depends on the measured per-task overhead of the specific runtime, describe the recursive cutoff as the implementation, and end with "then I measure" — including measuring whether parallelising helped at all, since for small inputs the honest answer is that it does not.

  • How would you determine the right sequential cutoff for a specific workload rather than guessing?
    Measure the per-task overhead of the runtime with a microbenchmark of trivial tasks, then time the sequential routine across input sizes to find the size whose duration is roughly fifty to a hundred times that overhead. Check that the resulting unit count is still several times the worker count, then sweep the cutoff over a range and plot the curve — it is usually flat across an order of magnitude, so being close is enough.
  • Why does making units too coarse hurt even when the total work is identical?
    Because wall clock is set by the slowest worker and the tail is bounded by one unit's duration. With as many units as workers, a single unit that runs long leaves every other worker idle for that period, and once a unit has started it cannot be redistributed. Over-partitioning gives the scheduler material to rebalance with.
  • Does the right granularity change when tasks block on I/O rather than compute?
    Yes. For blocking tasks the per-unit overhead is small relative to a multi-millisecond wait, so much finer units are affordable, and the constraint shifts from scheduling cost to how many concurrent in-flight operations the downstream resource tolerates. The sizing question becomes concurrency limits and connection budgets rather than task overhead ratios.

Handing out work on a building site. Assigning one brick at a time means the foreman spends all day talking and nobody lays bricks; assigning one whole wall each means the person with the awkward wall finishes hours after everyone else. You hand out sections — big enough to be worth the walk, small enough that nobody is stuck alone at the end.

saying these in an interview costs you the question

  • Quoting a fixed universal task size with no reference to the runtime's measured per-task overhead.
  • Creating one task per element for elements that take nanoseconds, so the envelope costs more than the work.
  • Creating exactly one unit per worker and assuming perfect balance.
  • Assuming finer granularity always improves load balance without accounting for the linear growth in overhead and queue contention.
  • Never concluding that the input is simply too small to parallelise.

context