In external merge sort, why does the pass count decide runtime rather than the comparison count?
answer
- count bytes moved, not comparisons
- what one pass costs at 500 GB
- nanoseconds against milliseconds
- the pass count is an integer staircase
- fan-in and buffer size compete for memory
basics
~20 sA pass reads and writes every byte, so one pass over 500 GB moves about 1 TB, taking minutes of device time; comparisons cost nanoseconds and hide under the I/O. Three passes instead of two is fifty percent more work.
solid answer
~50 sBoth a two-pass and a three-pass plan do Θ(n log n) comparisons, so the comparison count cannot distinguish them — yet one is fifty percent slower. The reason is the cost model: the scarce resource is bytes crossing the memory/storage boundary, and one pass moves `2N` bytes. At 500 GB that is 1 TB per pass, while comparisons hide entirely under that wait if you prefetch the next block while merging the current one. So the optimisation target is the *integer* pass count, plus the bytes moved within each pass. That also bounds how far fan-in should be pushed: with fixed memory, raising fan-in shrinks each stream's buffer, and once buffers get small the merge issues many small scattered reads and loses sequential bandwidth. Take the fan-in you need to hit the pass count, then spend all remaining memory on bigger buffers.
go deeper
Remember that one pass reads and writes every byte, so a pass over 500 GB moves about 1 TB. That single fact explains why fewer passes matters far more than faster comparisons.
Explain why two plans with identical n log n comparison counts can differ by fifty percent in runtime, and show the conversion from pass count to bytes moved to an estimated elapsed time at device bandwidth.
Demonstrate the fan-in versus buffer-size tension under fixed memory, and give the diagnostic: compare achieved throughput to the device's sequential ceiling before touching anything else.
Own the tuning policy. Decide which quantities are worth engineering effort at all — pass count and bytes per pass — and set the expectation that work not crossing a pass boundary is not funded.
## Two cost models, and why the usual one is the wrong one The habitual model for sorting counts comparisons, and it says every comparison sort is Ω(n log n) and the good ones are Θ(n log n). That model is fine when the data is in memory, because a comparison and a data access cost about the same. External sorting breaks the assumption: a comparison costs on the order of a nanosecond, while pulling a block off durable storage costs microseconds to milliseconds. When one operation is thousands to millions of times more expensive than the other, an accounting that counts the cheap one is not a model — it is noise. The replacement model counts **passes**. One pass reads every byte once and writes every byte once, so a pass over `N` bytes moves `2N`. For a 500 GB clickstream archive that is roughly 1 TB per pass. A two-pass sort moves about 2 TB; a three-pass sort moves about 3 TB. Both do Θ(n log n) comparisons. One is fifty percent slower. This is exactly the situation the wrong answer misses: *"big-O says n log n either way, so the pass count is a detail."* The pass count is not a detail; it is the answer. ## Why the comparisons genuinely disappear It is worth being precise about why the CPU side can be ignored rather than merely asserted. During a merge the process is doing two things concurrently: comparing the current heads of the open runs, and waiting for storage. With double buffering — while block `k` of a run is being merged, block `k+1` is already in flight — the comparison work runs inside the I/O wait. The merge then proceeds at device bandwidth, and the CPU has slack. Make the comparator twice as fast and nothing changes; the job still waits on storage. Remove one pass and a third of the elapsed time disappears. The practical corollary is that optimisation effort in an external sort should be aimed at exactly two quantities: - **The pass count** — the integer number of times the data crosses the boundary. - **The bytes moved per pass** — reduce the records themselves and every pass gets cheaper: compress runs, project away columns nobody needs downstream, or sort a compact `(key, record-pointer)` pair instead of whole records when the payload is large. Anything else — a faster in-memory sort in phase one, a cheaper comparator, more threads on comparisons — is optimising the resource that is already idle. ## The step function Because passes are integers, the payoff curve is a staircase, not a slope. Suppose the run count sits at 1000 and the fan-in at 127: two merge passes, three total. Nudging fan-in to 200 changes nothing — still two merge passes, since 1000 is still above 200. Getting fan-in to 1000 collapses the merge to one pass and removes a third of the total I/O in one step. So the useful question about any proposed tuning is not "does this improve a metric?" but "does this cross a boundary?" — the run count passing under the fan-in, or under fan-in squared. Improvements that do not cross a boundary buy nothing measurable. ## The tension nobody states until you are senior With fixed sort memory `M` and buffer size `b`, the fan-in is `F = M/b - 1`. Fan-in and buffer size trade directly against each other, and both matter: - **Shrink `b`** and `F` rises, which lowers the pass count. But each read of each stream becomes small. With many streams active, the device sees interleaved small requests across `F` different regions — the request pattern drifts from sequential toward scattered, per-request latency stops being amortised over a big transfer, and effective bandwidth collapses. You can win the pass count and lose more than you won. - **Grow `b`** and each stream is read in long sequential chunks at full bandwidth. But `F` falls, and if it falls below the run count you buy an entire extra pass. The rule that follows is an ordering, not a formula: first find the smallest fan-in that achieves your target pass count, then spend every remaining byte of memory on making the buffers as large as possible. Maximising fan-in for its own sake is a classic self-inflicted wound — the job hits its two-pass target and still runs at a quarter of device bandwidth because each of a thousand streams is being read 64 KB at a time. ## Diagnosing the gap This gives a concrete diagnostic when an external sort runs far slower than `passes × 2N / bandwidth` predicts. Compare achieved throughput against the device's sequential ceiling. If the job is at a small fraction of it, the merge is not streaming: check whether fan-in was pushed so high that per-stream buffers went tiny, whether read-ahead is actually happening or every block is fetched on demand with the CPU idle in between, and whether the runs and the output share a device that is now servicing reads and writes by turns. If throughput is near the ceiling and the job is still too slow, the pass count or the bytes per pass is the only thing left to attack — and at that point more memory, fewer columns, or compressed runs are the real levers.
- You have memory to spare — is a larger merge fan-in always better?No. Once the fan-in is high enough to finish the merge in the target number of passes, extra fan-in only shrinks each stream's buffer. Many streams read in small blocks turn a sequential workload into scattered requests and can cost more bandwidth than the saved pass was worth. Spend spare memory on bigger buffers, not on more open runs.
- Where does the comparison work go if I/O dominates?Under the I/O. With double buffering, the next block of each run is being fetched while the current one is merged, so comparisons run inside the wait and the sort proceeds at device bandwidth. That is why a faster comparator changes nothing measurable while removing a pass changes everything.
- An external sort runs three times slower than the pass arithmetic predicts — where do you look?Compare achieved throughput to the device's sequential ceiling. If it is far below, the merge is not streaming: fan-in pushed so high that per-stream buffers are tiny, no effective read-ahead so every block is fetched on demand, or runs and output contending on the same device. If throughput is near the ceiling, only the pass count or bytes per pass are left to attack.
- Which optimisations actually reduce the bytes moved per pass?Compressing the runs, projecting away fields no downstream consumer reads, and sorting compact key-plus-pointer pairs instead of whole records when payloads are large. Each shrinks the data that every pass must carry, so the saving multiplies by the pass count.
saying these in an interview costs you the question
- Says n log n comparisons means the plan is already optimal
- Counts a pass as reading the data only, not reading and writing
- Maximises fan-in without checking the resulting buffer size
- Assumes sequential and scattered device throughput are comparable
- Optimises the in-memory sort while the merge still takes three passes
- Treats the pass count as continuous rather than an integer staircase