In a recursive divide-and-conquer parallel algorithm, why do implementations stop splitting once a subproblem is below some size and run the remainder sequentially, and how would you choose that threshold?
answer
- fixed cost per task vs geometrically shrinking work
- W/(W+C): stop when W approaches C
- aim for ~10+ leaves per worker
- plateau in a cutoff sweep, not a single best point
- depth-based cutoff adapts to core count
basics
~20 sEach split costs bookkeeping - creating, queueing and scheduling a task. Below some size that overhead exceeds the useful work, so the recursion would slow things down. Pick the threshold by measuring: choose the smallest size where parallel still beats sequential, comfortably above the break-even point.
solid answer
~60 sSplitting is not free: every fork allocates a task, pushes it on a queue, and may be picked up by another worker, costing cache misses and scheduling work. The useful work per leaf shrinks geometrically as you recurse, so eventually the overhead per task exceeds the computation in it and the parallel version becomes slower than a plain loop. So you set a **sequential cutoff**: below N elements, run the ordinary sequential algorithm. Choosing N is empirical. A common rule is to target enough leaves for good load balance - on the order of a few times the number of workers, often ten-plus tasks per worker so a stolen task can rebalance - while keeping each leaf big enough that a task's overhead is a small percentage of its runtime. In practice: measure the sequential cost of one element's work, measure task overhead, and pick a cutoff where per-leaf work is one to two orders of magnitude larger. Then benchmark on the real input and hardware, because the right number depends on both.
code
text · 10 lines// size cutoff
if length(range) <= CUTOFF: return sequential(range)
// depth cutoff: stop after enough tasks exist
solve(range, depth):
if depth == 0 or length(range) <= MIN:
return sequential(range)
... fork solve(a, depth-1); solve(b, depth-1) ...
// depth chosen so 2^depth ~= k * workerCount, k ~ 8..32go deeper
Know that splitting has a per-task cost, so tiny subproblems are run with a normal sequential loop instead of being split further.
Explain the two opposing forces - task overhead versus load balance - and give a rough target such as many more leaves than workers with each leaf's overhead well under a few percent.
Describe how you would actually measure it: per-element cost, per-task overhead, a cutoff sweep on real data, and picking the middle of the plateau so the choice survives a different machine.
Discuss making granularity adaptive rather than constant - depth-based or idleness-driven cutoffs - and the operational risk of a hardcoded threshold across heterogeneous hardware and varying per-element cost.
## Why an unbounded recursion is wrong A divide-and-conquer parallel algorithm could in principle recurse until each leaf holds a single element. It does not, because the model has a fixed cost per task that has nothing to do with the size of the task. Each fork must: allocate or reuse a task record, publish it where a worker can find it, be scheduled, possibly be migrated to another core (dragging its data through a cold cache), and finally be joined, with the parent's continuation resumed. Call this cost **C**. If a leaf performs work **W**, the leaf's efficiency is roughly W / (W + C). When W falls to the same order as C you spend half your machine on bookkeeping; when W falls below C you are strictly slower than a sequential loop - and you have burned memory and cache lines to be slower. Because the recursion halves the problem, W shrinks geometrically while C stays constant. Two extra levels of recursion quadruple the number of leaves and quarter the work in each. That is why the last few levels of an unbounded recursion contain almost all the tasks and almost none of the value. ## The cutoff The fix is a base case that is a *size*, not a size of one: ``` solve(range): if length(range) <= CUTOFF: return solveSequentially(range) // plain loop, no tasks ...split, fork, join, combine... ``` The leaf now uses the ordinary sequential algorithm, which is also usually the *fastest* code you have for small inputs: it keeps everything in registers and cache, has no allocation, and may be vectorized by the compiler. ## Two forces pulling in opposite directions **Push the cutoff up (fewer, bigger tasks):** less overhead, better cache locality, less memory churn. **Push the cutoff down (more, smaller tasks):** better load balance. If your tasks are few and one of them is slow - because the data is uneven, or a core is busy with something else - the whole computation waits for that straggler. Many small tasks let idle workers pick up spare work and smooth the tail. This is why a rule of thumb is to aim for on the order of ten or more leaves per worker rather than exactly one per worker: with one task per core, a single slow core doubles your finish time. So the cutoff is a **granularity** decision, and the target is: small enough that there are plenty of leaves to rebalance, large enough that per-task overhead is negligible. ## A concrete way to pick it 1. Measure the sequential cost of the leaf operation per element - say 5 ns per element for a sum, 100 ns per element for something heavier. 2. Estimate the per-task overhead of your runtime - typically hundreds of nanoseconds to a few microseconds including the steal and cache effects. 3. Choose leaf work one to two orders of magnitude above the overhead. If a task costs ~1 microsecond of overhead and you want overhead under 1 percent, target ~100 microseconds of work per leaf. At 5 ns per element that is about 20,000 elements. 4. Sanity check leaf count: total_size / CUTOFF should be comfortably more than the worker count. If it is not, lower the cutoff even at some overhead cost - load imbalance hurts more. 5. Benchmark a sweep of cutoffs on realistic input and hardware, and look for the plateau rather than the single best point. The curve is typically flat over a wide range and falls off sharply at the small end; pick the middle of the plateau so you are robust to a different machine. ## Adaptive alternatives Some runtimes avoid a fixed element count by cutting on **depth** instead: split until you have produced roughly k times the worker count of tasks, then go sequential. This adapts automatically to machine size but not to per-element cost. Others cut adaptively: keep splitting only while other workers appear to be idle, so a saturated machine stops producing tasks. Both are refinements of the same idea - stop when extra parallelism has no one to run on. ## Common mistakes Hardcoding a cutoff tuned on a laptop and shipping it to a 64-core server; expressing the cutoff in elements when the per-element cost varies by orders of magnitude between call sites; or setting it so high that you produce four tasks on a sixteen-core machine and then wondering why speedup stalls near four. The honest interview answer is that there is no universal number: it is per-algorithm, per-machine, and it must be measured. What an interviewer wants to hear is *why* the number exists and which two forces it balances.
- Why not simply set the cutoff to total_size / numberOfWorkers, so each worker gets exactly one task?Because that assumes every task takes the same time and every worker is fully available. Any straggler - uneven data, a core shared with another process, an unlucky cache profile - then extends the whole computation, because there is no spare task for an idle worker to pick up. Oversubscribing with roughly an order of magnitude more tasks than workers lets the runtime rebalance and keeps the tail short.
- How does a cutoff expressed in element count go wrong when the per-element work varies?The cutoff is really about work per leaf, and element count is only a proxy for it. If one call site does 5 ns per element and another decodes an image per element, the same element threshold gives leaves that differ by orders of magnitude in duration - too much overhead in one case, too coarse and unbalanced in the other. Either tune per call site or express the cutoff in estimated work rather than raw count.
Delegating work by memo: sending a memo costs the same whether the job takes an hour or ten seconds. Past some point you just do the small jobs yourself.
saying these in an interview costs you the question
- Claiming the cutoff is a fixed universal number such as 1000, valid everywhere.
- Recursing to single elements because 'more parallelism is always better'.
- Setting the cutoff so high that fewer tasks exist than there are cores.
- Tuning the threshold once on one machine and treating it as a property of the algorithm.
- Ignoring load balance entirely and optimizing only for the lowest task count.