skip to content

You are splitting a large collection across worker threads. Compare contiguous-block partitioning, round-robin partitioning, and dynamic chunk handout, and say which you would choose when the per-item cost varies a lot.

level: middleimportance: should knowfreq 40%

answer

  1. block = locality, no balance
  2. cyclic = balance for smooth gradients, kills locality
  3. dynamic = balance regardless of distribution, pays an atomic per chunk
  4. over-partition 8-16x P; partition by cost, not count
  5. one giant item = a floor no scheme beats

basics

~20 s

Block gives each worker one contiguous range — best locality, worst load balance under uneven costs. Round-robin interleaves items across workers — balances a smooth cost gradient, hurts locality. Dynamic handout gives workers a new chunk when they finish — best balance under unpredictable costs, at the price of coordination on the shared cursor. With highly variable per-item cost, use dynamic (or many small static chunks).

solid answer

~60 s

**Contiguous block**: worker *k* gets items [k·N/P, (k+1)·N/P). Zero coordination, perfect cache and prefetch behaviour, ideal for uniform per-item cost. Its failure mode is skew: if the expensive items cluster in one range, one worker runs long and everyone else idles. **Round-robin / cyclic**: worker *k* gets items k, k+P, k+2P… Statistically evens out a cost that varies smoothly with position (e.g. rows that get progressively longer). Costs you spatial locality — each worker touches scattered addresses, and adjacent items handled by different workers can cause false sharing on the output. **Dynamic handout**: a shared cursor; each worker grabs the next chunk when it finishes. Self-balancing regardless of the cost distribution — the tail is bounded by one chunk, not one range. Cost: atomic contention on the cursor, worse locality, non-deterministic assignment. With highly variable cost, choose dynamic — or the practical middle ground, **static over-partitioning**: cut the input into far more chunks than workers (say 8–16× P) and hand them out. You get most of the balance with far less cursor traffic. If you can *estimate* per-item cost, partition by estimated cost rather than count and keep the cheap block strategy.

code

text · 12 lines
text
items:      0 1 2 3 4 5 6 7 8 9 10 11

BLOCK:      W0: 0..3     W1: 4..7     W2: 8..11
            contiguous, no coordination, bad if 8..11 are the expensive ones

CYCLIC:     W0: 0,3,6,9  W1: 1,4,7,10 W2: 2,5,8,11
            interleaved, evens out a rising cost curve, scattered memory access

DYNAMIC (chunk = 2, shared cursor):
            W0 takes 0-1, W1 takes 2-3, W2 takes 4-5,
            whoever finishes first takes 6-7, and so on
            final imbalance <= one chunk, cost = one atomic claim per chunk

go deeper

for a junior

Know that work can be split into contiguous ranges or handed out on demand, and that uneven work leaves some workers idle.

for a middle

Compare the three schemes on balance, locality and coordination, and pick dynamic or over-partitioned static when per-item cost is unpredictable.

for a senior

Bring in chunk-size selection, guided tapering, false sharing on shared output, cost-estimate partitioning, and the hard floor set by a single oversized item.

for a principal

Discuss skew as a systemic problem — hot keys, cost models, over-partitioning as a policy — and how the partitioning choice interacts with the scheduler and with reproducibility requirements.

## The problem partitioning solves Given N items and P workers, you must decide which worker processes which items. The decision trades three things against each other: 1. **Load balance** — is every worker busy until the end? Wall clock equals the *slowest* worker, so an imbalance of 30% wastes 30% of your machine. 2. **Locality** — do a worker's items sit near each other in memory, so hardware prefetching and cache lines work for you? 3. **Coordination overhead** — how much do workers pay to find out what to do next? No scheme wins all three; the workload's cost distribution decides which one to sacrifice. ## Contiguous block (static, range) Worker *k* takes a single contiguous range of roughly N/P items. - **Coordination**: none. Ranges are computed once, arithmetically. - **Locality**: excellent. Each worker sweeps a contiguous region — prefetchers work perfectly, and on NUMA hardware the range can be allocated on the worker's own memory node. - **Balance**: only if per-item cost is uniform. The failure mode: the last worker gets all the long strings, or the image tiles at the bottom of a fractal all take 100× longer. Wall clock becomes the cost of the worst range, and utilisation can collapse to a fraction of the machine. This is the correct default for uniform numeric work — sums over arrays, per-pixel transforms, fixed-size record parsing. ## Round-robin / cyclic (static, interleaved) Worker *k* takes items k, k+P, k+2P, … - **Coordination**: none — still computed arithmetically. - **Balance**: good when cost varies *smoothly with position*. The classic case is a triangular loop where iteration *i* does *i* units of work: block partitioning gives the last worker the heaviest half, while cyclic gives everyone a fair mix of cheap and expensive. - **Locality**: poor. Each worker touches every P-th element, so cache lines are shared between workers and prefetching is defeated. If workers *write* to adjacent output elements, you also get **false sharing** — different workers writing different variables that happen to sit on the same cache line, forcing the line to bounce between cores and destroying performance even though the code is logically correct. **Block-cyclic** is the standard compromise: hand out contiguous chunks of size B in round-robin fashion. Chunk size B tunes between the two — larger B recovers locality, smaller B improves balance. ## Dynamic handout (self-scheduling) A shared cursor (an atomic counter) over the input. Each worker atomically claims the next chunk, processes it, and comes back for more. - **Balance**: excellent and *distribution-independent*. A worker that draws an expensive chunk simply claims fewer chunks. Worst-case imbalance at the end is one chunk's duration, so it is bounded by chunk size rather than by range size. - **Coordination**: an atomic operation per chunk. With chunk size 1 and cheap items this dominates — the counter becomes a contention point and can make the parallel version slower than sequential. - **Locality**: worse than block, better than cyclic if chunks are contiguous. - **Determinism**: assignment varies run to run, which complicates reproducibility and debugging. **Guided / tapered chunking** is the refined form: hand out large chunks early (cheap coordination, good locality while there is plenty of work) and progressively smaller chunks near the end (fine-grained balancing exactly where the tail matters). A common rule is to hand out about 1/P of the *remaining* work each time. ## Choosing under variable cost With high per-item cost variance: 1. **If cost is unpredictable** — dynamic handout or guided chunking. This is the only scheme that adapts to a distribution you cannot see in advance. 2. **If cost is predictable per item** (you can estimate from a field: file size, string length, row count) — partition by *estimated cost*, not by count. Sort or bucket items so each partition has roughly equal predicted work, then use cheap static block assignment. This gets balance and locality together. 3. **If you cannot do either** — static over-partitioning: cut into 8–16× P chunks and hand them out. It gives most of the balance benefit while amortising the cursor cost over a chunk, and it is the default in most parallel frameworks (which then use work stealing to redistribute). ## Skew is the enemy, and it has a floor If a single item costs more than total_work / P, no partitioning scheme can beat that item's duration — it cannot be split across workers. The only remedies are splitting the item itself (if it is divisible), handling it with a specialised path, or accepting the tail. In distributed data processing this is the "hot key" problem, and the standard mitigations are the same: salt the key to split it, pre-aggregate before the shuffle, or route the heavy key to a dedicated path. ## Practical checklist - Uniform cost, memory-bound: contiguous block, one chunk per worker. - Cost varying with position: block-cyclic with a moderate chunk size. - Unknown or bursty cost: dynamic or guided chunking, chunk sized so each chunk is far larger than the claim overhead. - Predictable cost: partition by cost estimate, keep blocks. - Writing to shared output arrays: pad or block-align per-worker regions so two workers never write the same cache line.

  • How do you pick the chunk size for dynamic handout?
    Large enough that the per-chunk claim cost (one atomic operation plus cache effects) is a negligible fraction of the chunk's work — a common target is well under one percent — and small enough that the final imbalance, which is bounded by one chunk, stays acceptable. Guided chunking sidesteps the tradeoff by starting with big chunks and shrinking them as remaining work runs out.
  • One item in the input costs more than a third of the total work, and you have three workers. What is the best possible wall clock?
    At least the duration of that single item, because it cannot be split across workers. No partitioning strategy improves on that floor. Your options are to subdivide the item itself if it is internally divisible, give it a specialised faster path, start it first so it overlaps everything else, or accept the tail and size expectations around it.
  • Why can round-robin partitioning make a write-heavy loop dramatically slower?
    Because adjacent output elements handled by different workers usually share a cache line. Each write invalidates the line in the other cores' caches, so the line ping-pongs between them — false sharing. The code is correct but memory traffic explodes. Contiguous blocks, or per-worker output regions padded to cache-line boundaries, avoid it.

saying these in an interview costs you the question

  • Always splitting into exactly one chunk per worker, which guarantees a long tail whenever costs are uneven.
  • Using per-item dynamic handout for cheap items, so the shared cursor's atomic traffic dominates the actual work.
  • Ignoring locality and false sharing when choosing an interleaved partitioning.
  • Assuming equal item counts mean equal work when per-item cost varies by orders of magnitude.
  • Believing dynamic scheduling can fix a single oversized item; it cannot be split across workers.

context