skip to content

Why does sorting 500 GB of clickstream events on a 4 GB machine need an external merge sort?

level: juniorimportance: must knowfreq 55%

answer

  1. the data outruns memory, not the CPU
  2. what any in-memory sort assumes about access
  3. random seeks versus streaming bandwidth
  4. sort what fits, write it back sorted
  5. make runs first, then merge them

basics

~20 s

External merge sort is needed because 500 GB will not fit in 4 GB, and an ordinary sort assumes free random access. It instead sorts memory-sized chunks into runs on disk, then merges them sequentially.

solid answer

~50 s

The problem is not the comparison count, it is the access pattern. Any in-memory sort jumps around the whole array — partitioning, merging and sifting all assume random access is free — and when the array is 125 times larger than memory, almost every one of those accesses becomes a disk seek. External merge sort restructures the work into two phases so that every access is sequential and bounded: **phase one** reads the input in memory-sized chunks, sorts each chunk in memory, and writes it back as a sorted *run*; **phase two** merges those runs by streaming a block from each and repeatedly emitting the smallest head. Both phases read and write in large sequential blocks, so the drive delivers streaming bandwidth instead of seek latency. Comparison work is unchanged at Θ(n log n); what changed is that the algorithm now controls when data crosses the memory boundary.

go deeper

for a junior

Be ready to say plainly that 500 GB cannot be held in 4 GB, and to name the two phases: sort memory-sized chunks into sorted runs, then merge the runs. Knowing that both phases read and write sequentially is enough at this level.

for a middle

Explain the mechanism, not just the shape: why an in-memory sort's random access pattern becomes disk seeks, why paging cannot rescue it, and why the run size is set by available memory rather than by the input size.

for a senior

Expect to be asked what you would actually run on a real box: how much memory you would give the sort buffer, how you would size the output buffer, and how you would confirm from device metrics that the job is streaming rather than seeking.

for a principal

Own the framing decision. Argue when sorting is the right shape at all versus partitioning by key, indexing, or pushing ordering to the producer — and when a job's ordering requirement should be renegotiated rather than engineered around.

## The constraint that changes the algorithm Sorting 500 GB of clickstream events on a machine with 4 GB of usable memory is not a harder *comparison* problem than sorting 4 GB. It is the same comparison problem under a different cost model. Once the data no longer fits in memory, the expensive resource stops being CPU cycles spent comparing and becomes **bytes moved between memory and durable storage**, and the gap is not a constant factor: a comparison costs on the order of nanoseconds, while a random read from a spinning disk costs milliseconds — a factor of roughly a million — and even on solid-state storage a scattered 4 KB read delivers a small fraction of the sequential bandwidth of the same device. ## Why an in-memory sort fails, mechanically Every classic in-memory sort is built on the assumption that indexing any position costs the same as indexing the neighbouring one. Partition-based sorts sweep two cursors from opposite ends of a range. Merge-based sorts read two sub-ranges and write a third. Heap-based sorts sift a value down a tree whose children live at positions `2i+1` and `2i+2`, which for large `i` are far apart. Feed any of them an array 125 times larger than memory and the machine cannot hold the working set; each of those "free" indexing operations becomes a request to storage, and because the pattern has no locality, the requests are effectively random. The common wrong answer is to keep the in-memory sort and delegate the problem downward: memory-map the file, or add enough swap, and let the operating system's paging handle it. The mapping itself succeeds — a large address space can map a file far bigger than memory — but the paging layer can only react to the access pattern it is given. It cannot know that the sort will want a particular page again in ten million comparisons, so it evicts pages the sort still needs, faults them straight back in, and the run degenerates into thrashing. The operating system is not sorting; it is servicing random faults for an algorithm that was designed on the assumption they would never happen. ## The two-phase restructuring External merge sort keeps the comparison work but rearranges *when* data crosses the memory boundary. **Phase one — run generation.** Read as much of the input as fits in the sort buffer, sort that chunk entirely in memory with any ordinary sort, and write it back to storage as a sorted *run*. Repeat until the input is consumed. The reads are sequential (you are scanning the input front to back), the writes are sequential (you are appending a sorted block), and the run size is bounded by memory, not by input size. With 500 GB and roughly 4 GB usable, this produces about 125 runs. **Phase two — merging.** Open several runs at once, keep one input buffer per run plus one output buffer, and repeatedly emit the smallest of the current heads, refilling a buffer from storage whenever it drains. Each run is consumed strictly front to back, so every read is a large sequential block; the output is written in large sequential blocks too. If more runs exist than can be merged at once, merge them in groups and repeat on the resulting longer runs. ## The cost model that replaces comparisons The standard accounting unit is the **pass**: one pass reads every byte once and writes every byte once. Run generation is one pass. Each round of merging is another. At 500 GB, one pass moves about 1 TB. This is why practitioners state external sort costs as "two passes" or "three passes" rather than in comparisons — comparisons are hidden under the I/O and can be overlapped with it by prefetching the next block while the current one is being merged. ## What does not change Two things are worth stating explicitly, because both are common misconceptions. First, the total number of comparisons is still Θ(n log n) — external merge sort is not a worse algorithm, it is the same algorithm with an I/O-aware schedule. Second, external sorting is not tied to any particular in-memory sort: phase one can use whichever sort is fastest on a memory-sized chunk, and stability of the whole result follows if that sort is stable and the merge breaks ties in favour of the earlier run. ## Where it shows up Anything that must produce ordered output over data larger than working memory ends up here: sorting a query result set that exceeds the memory a sort operator was given, ordering a day of clickstream events by session and timestamp before sessionization, or producing sorted input for a merge-based join. The recognisable signal in an interview is the phrase "does not fit in memory" — it is an instruction to switch cost models, not to look for a cleverer comparison sort.

  • Why doesn't memory-mapping the file and letting the operating system page it in solve this?
    Because paging is reactive. The sort's accesses are scattered across the whole 500 GB range with no locality, so the paging layer evicts pages that will be needed again and faults them back in — the sort becomes random disk I/O with extra overhead. External merge sort works because it *chooses* the access order in advance and makes it sequential, which no paging heuristic can infer on its behalf.
  • How large should each sorted run in phase one be?
    As large as the sort buffer allows — usable memory minus whatever the process needs for output buffering and overhead. Larger runs mean fewer runs, and fewer runs mean a better chance of finishing the merge in a single pass. The run count is simply the input size divided by the run size, rounded up.
  • Does external merge sort perform more comparisons than an in-memory sort?
    No — it is still Θ(n log n) comparisons overall, and the merge phase does no redundant comparison work. The extra cost is entirely I/O: the data is read and written once per pass. That is why the analysis is stated in passes and bytes moved rather than in comparison counts.

You cannot alphabetise a warehouse of paper records on one small desk. You sort a deskful at a time into ordered stacks, then walk the stack tops in order, taking the smallest each time.

saying these in an interview costs you the question

  • Says just add swap and let the operating system handle it
  • Claims external sorting needs more comparisons than an in-memory sort
  • Treats external sorting as a different comparison algorithm, not an I/O strategy
  • Assumes random and sequential storage access cost about the same
  • Thinks the run size should be chosen from the input size rather than from memory

context