A recursive parallel sort splits an array in half, sorts both halves concurrently, and then merges the two sorted halves with an ordinary sequential merge. Even on a machine with many idle cores the speedup flattens out. How do you reason about the ceiling, and how would you restructure the decomposition to raise it?
answer
- work T1 vs span T-inf; parallelism = T1 / T-inf
- combine step sets the ceiling
- sequential merge -> span O(n) -> speedup ~ log n
- parallel merge by median + binary search -> span O(log^2 n)
- idle cores + flat speedup = span-bound, not scheduler-bound
basics
~20 sThe sequential merges form a chain: the final merge alone touches every element, and nothing can shorten it. That chain is the critical path of the task graph and bounds the finish time no matter how many cores exist. Raise the ceiling by parallelizing the merge itself.
solid answer
~60 sReason with two quantities from the task graph: **work** (total operations, what one worker would do) and **span** (the longest chain of dependent operations, the finish time with unlimited workers). Maximum speedup is at most work divided by span. For this sort, work is the usual n log n. But the top-level merge is sequential and costs about n, the merge below it about n/2, and so on - a chain of roughly n + n/2 + n/4 ... which is linear in n. So the span is dominated by the merges and speedup saturates near log n, not near the core count. Doubling cores stops helping long before you run out of them. The fix is to attack the span, not the work: make the merge itself divide-and-conquer - binary-search the median of one run in the other, split both runs at that point, and merge the two pairs concurrently. That reduces merge span to polylogarithmic and lifts the ceiling by orders of magnitude, at the cost of more work and complexity.
code
text · 10 linesparallel sum: S(n) = S(n/2) + O(1) -> S = O(log n)
W(n) = 2W(n/2) + O(1) -> W = O(n)
parallelism = n / log n (scales well)
sort, seq merge: S(n) = S(n/2) + O(n) -> S = O(n)
W(n) = 2W(n/2) + O(n) -> W = O(n log n)
parallelism = log n (saturates early)
sort, par merge: S(n) = S(n/2) + O(log^2 n) -> S = O(log^3 n)
parallelism ~ n / log^2 n (machine-bound)go deeper
Recognize that the final merge is sequential work that every element must pass through, so extra cores cannot speed it up.
Use the work-versus-span vocabulary: the longest dependency chain sets the best possible finish time, and here the chain of merges is linear in n.
Derive the recurrences, conclude that parallelism is only about log n, and explain diagnosis - flat speedup with idle workers means span-bound, so tuning cutoffs will not help.
Treat scalability as a decomposition-time design choice: compare parallel-merge versus partition-based algorithms, weigh added work and complexity against ceiling gained, and set expectations for the target hardware before committing.
## Reasoning tool: work and span A fork-join computation is a directed acyclic graph of operations: forks branch it, joins rejoin it. Two numbers characterize it. - **Work (T1)**: the total number of operations - the time on one worker. - **Span (T-infinity)**, also called depth or critical path: the length of the longest chain of operations where each must finish before the next begins - the time on infinitely many workers. The best possible speedup is **T1 / T-infinity**, the *parallelism* of the algorithm. It is a property of the decomposition, not of the machine. If parallelism is 8, buying a 64-core machine buys you nothing past 8. This is why 'the cores are idle' is not evidence that the runtime is at fault - it may simply be that the graph has no more independent work at that moment. ## Applying it to the sort Split the array, sort halves concurrently, merge sequentially. - **Work**: the recurrence is W(n) = 2 W(n/2) + O(n), giving the familiar O(n log n). Parallelizing does not change the work, and should not - an algorithm that does far more total work to gain parallelism often loses on real machines. - **Span**: the two recursive sorts run concurrently, so only one contributes; then the merge is added. S(n) = S(n/2) + O(n). That expands to n + n/2 + n/4 + ... = O(n). The merge chain is linear. So parallelism = O(n log n) / O(n) = O(log n). For a million elements that is about 20 - meaning speedup saturates around twenty-fold no matter how many cores you own, and in practice much lower once overheads count. The last merge alone reads and writes every element with no internal parallelism; it is a hard floor. ## Where the intuition usually goes wrong People notice that the *sorting* is fully parallel and conclude the whole thing is. But a chain is only as short as its longest link, and the combine steps of a divide-and-conquer algorithm sit on the critical path exactly like the split steps do. Whenever the combine is sequential and linear in the subproblem size, it dominates the span. This generalizes into a design rule: **in divide-and-conquer, the combine step determines your scalability ceiling.** Parallel sum has an O(1) combine (one addition), so its span is O(log n) and its parallelism is O(n / log n) - excellent. Parallel sort has an O(n) combine, and pays for it. ## Restructuring: parallelize the merge The standard fix makes merging itself divide-and-conquer: ``` parallelMerge(A, B, out): if small: sequentialMerge(A, B, out); return pick the median element m of the longer run, at index i j = binarySearch(m, shorter run) // O(log n) place m at position i + j in out fork parallelMerge(A[..i], B[..j], out[..i+j]) parallelMerge(A[i+1..], B[j..], out[i+j+1..]) join ``` Each level does an O(log n) binary search and splits both inputs, so merge span becomes O(log^2 n) instead of O(n). Plugging that into S(n) = S(n/2) + O(log^2 n) gives a total span of O(log^3 n) and parallelism close to n / log^2 n - now the machine, not the algorithm, is the limit. The price: extra work for the searches, an out-of-place output buffer, and considerably harder code. ## Other decomposition levers - **Change the algorithm, not just the merge.** A sample-sort or bucket-based approach partitions once by value ranges, sorts buckets fully independently, and then only concatenates - a trivial combine. It moves the hard part into choosing good splitters, which can be done by sampling. - **Overlap combine with compute.** Merge pairs as soon as both are ready rather than waiting for a whole level, so the graph is not artificially level-synchronized. - **Reduce constants where the span lives.** The final merge is memory-bandwidth bound; a cache-friendly layout speeds up the exact step that dominates. ## How to decide in practice Do the span analysis before writing code: if the combine is O(size), expect logarithmic parallelism and do not promise linear scaling. Then measure: plot speedup against worker count and find where the curve flattens. If it flattens far below core count while CPU is idle, you are span-bound and no amount of scheduler tuning will help - the decomposition must change. If instead workers are busy but speedup is poor, you are overhead- or memory-bound, a different problem with different fixes. The principal-level point is that scalability is designed in at decomposition time. Choosing an algorithm whose combine step is cheap is worth far more than any amount of tuning applied afterwards.
- Two decompositions have the same span but one does 40 percent more total work. Which do you pick?Usually the lower-work one, unless the target machine has enough spare cores to absorb the extra work and the span advantage translates into real wall-clock gain. Extra work is paid on every run and consumes memory bandwidth and energy, while extra parallelism only pays off when cores are actually idle. The decision is workload- and hardware-specific, so it is settled by measuring speedup curves, not by the asymptotics alone.
- You measure speedup flattening at 6x on a 32-core machine. How do you tell whether you are span-bound or overhead-bound?Look at worker utilization. Span-bound looks like idle workers with nothing to steal - the graph has no independent work available at that moment - and the flattening point does not move when you change the cutoff. Overhead-bound looks like busy workers with poor efficiency, and it responds to a coarser cutoff, fewer allocations, or better locality. Instrumenting the number of ready tasks over time separates the two quickly.
- Why does a sample-sort style partition raise the ceiling more cheaply than parallelizing the merge?Because it removes the expensive combine altogether: once elements are partitioned into value-ordered buckets, the sorted buckets simply concatenate, so the combine is O(1) per bucket. The cost moves to choosing good splitters, which can be done from a small random sample and is cheap relative to the data. The risk shifts from span to imbalance - skewed splitters give uneven buckets - which is a load-balancing problem rather than a scalability ceiling.
A relay race where every runner can be replaced by a team, except the final leg, which one runner must run alone. Adding people never beats the time of that last leg.
saying these in an interview costs you the question
- Assuming that because both halves sort in parallel, the whole algorithm scales linearly with cores.
- Blaming the scheduler or thread pool when the ceiling is imposed by the algorithm's critical path.
- Trying to fix a span-bound algorithm by lowering the sequential cutoff or adding threads.
- Treating extra total work as free whenever it buys more parallelism.
- Ignoring that the top-level combine step touches every element and cannot be skipped.