Explain how a database sorts a result set far larger than the memory budget given to the sort operation, and describe the I/O cost of that algorithm.
answer
- Run generation: fill, sort, write sorted run
- Merge: heap over run heads, one buffer per run
- Runs ~ N/M; fan-in ~ M/buffer
- Cost ~ 2N per pass; passes grow logarithmically
- Ordered index or Top-N heap can avoid the sort
basics
~20 sIt uses an external merge sort: fill memory, sort that chunk, write it out as a sorted run, repeat until input is exhausted, then merge the runs by repeatedly taking the smallest head value across them. Each pass reads and writes the whole dataset, and more memory means fewer, larger runs and usually a single merge pass.
solid answer
~60 sWhen input exceeds the sort's memory budget, the operator switches from an in-memory sort to an **external merge sort**, in two phases. **Run generation:** read input until the budget is full, sort that batch in memory, write it to a temporary file as a sorted *run*, then start the next batch. With M bytes of memory and N bytes of data you produce roughly N/M runs. **Merge:** open several runs at once, keep one input buffer per run, and repeatedly emit the smallest head value across them - typically driven by a heap - refilling a buffer whenever it drains. Output streams to the parent operator, so the sort can produce rows without materializing the whole sorted result. If the number of runs exceeds how many can be merged at once (the merge fan-in, bounded by memory divided by buffer size), the engine merges in several passes, each pass reading and writing the entire dataset. Cost is therefore roughly `2 x N x number_of_passes` of I/O on temporary files, versus zero for an in-memory sort. That is why moderately more memory can convert a multi-pass sort into a single-pass one and produce a large speedup.
code
text · 4 linesphase 1: 16 runs x ~512MB, each internally sorted -> temp files
phase 2: fan-in 16 <= max fan-in -> single merge pass
temp I/O: write 8GB + read 8GB = ~16GB
in-memory equivalent: 0 temp I/Ogo deeper
Describe the two phases - make sorted chunks on disk, then merge them - and that this is slower than sorting in memory.
Give the cost model: runs equal data over memory, fan-in bounded by memory, each pass costs a full read plus write.
Discuss single-pass versus multi-pass regimes, streaming the final merge, temp-space capacity, and avoiding the sort via an ordered index or a Top-N heap.
Reason about temp storage as a provisioned resource, where scratch I/O lands relative to the log, and whether the right fix is memory, statistics, indexing or a different plan shape.
## The problem A sort is a blocking operator: it cannot emit the first row until it has seen the last input row, because any late row might sort first. If the input fits in the operator's memory budget, the engine sorts in RAM and streams results out. When it does not fit, the data must be staged on disk - but naive approaches (random-access sorting on disk) would be catastrophically slow, so engines use an algorithm designed for sequential I/O: the **external merge sort**. ## Phase 1 - run generation Read input rows into the memory budget until it is full. Sort that batch in memory with a standard comparison sort. Write it to a temporary file as a **sorted run** - a sequential write. Clear memory and repeat. With N bytes of input and M bytes of budget you get about N/M runs, each roughly M bytes and internally sorted. A common refinement, replacement selection, uses a heap to produce runs averaging around 2M, halving the run count for partially ordered input. Two important details: - Engines often sort **keys plus a pointer or a compact tuple** rather than full rows, so more entries fit per byte. - Temporary files are written sequentially, which is far cheaper per byte than random access, and are usually not logged in the write-ahead log because they are scratch data that recovery can discard. ## Phase 2 - merging Merging k sorted runs into one is a linear streaming operation: keep a small read buffer per run, maintain a heap of the current head value of each, repeatedly pop the minimum, emit it, and refill from that run's buffer. The **fan-in** k is limited by memory: each of the k runs needs a read buffer, plus one output buffer, so k is roughly `budget / buffer_size`. With a typical budget this is comfortably in the tens or hundreds. - If `runs <= k`, one merge pass suffices, and output can be streamed directly to the parent operator - no need to store the final sorted result at all. - If `runs > k`, the engine merges groups of k runs into larger intermediate runs, and repeats. Each such pass reads and rewrites the whole dataset. ## Cost model Let N be the data size in pages, M the memory budget in pages, and B the buffer size per run. - Number of initial runs: about N/M. - Merge fan-in: about M/B. - Number of passes: `1 + ceil(log_fanin(N/M))`. - I/O: each pass reads and writes N pages, so roughly `2N x passes` of temporary I/O (the final pass can skip the write if output is streamed). The crucial shape is the logarithm: passes grow very slowly with data size, so the practical world has essentially two cases - **fits in memory** (zero temporary I/O) and **one merge pass** (2N of temporary I/O). Multi-pass sorts require a truly extreme ratio of data to memory and generally indicate a badly undersized budget. That is why the payoff curve from more sort memory is stepped, not smooth: going from spilling to in-memory is a step change, and going from multi-pass to single-pass is another. Doubling memory within the same regime changes little. ## Practical consequences - **Sorts are unavoidable** for ordering, for merge joins, for some grouping and distinct strategies, and for index builds - so external sorting is core machinery, not an edge case. - **A sorted index can remove the sort entirely.** If an ordered index supplies rows already in the required order, the plan skips the sort and its memory and I/O disappear, which is usually a bigger win than tuning memory. - **Temporary space is a real resource.** A large sort can consume many gigabytes of scratch space; running the temp area out of disk fails the query, and sharing that device with the transaction log couples analytical spills to transactional latency. - **Bad row estimates cause surprise spills.** If the optimizer expects a thousand rows and gets ten million, it may have chosen a plan and a memory grant on the wrong scale. The fix there is statistics or the plan, not memory. - **Top-N is special.** Ordering with a small limit lets the engine keep only the best N rows in a bounded heap, avoiding a full sort and any spill.
- Doubling the sort memory budget for a query that already spills sometimes gives a huge speedup and sometimes almost none. Why?The cost is stepped, driven by the number of passes rather than by memory smoothly. If the extra memory pushes the operator from spilling to fully in-memory, temporary I/O drops to zero - a large win. If it merely turns a multi-pass merge into a single-pass merge, you save one full read and write of the data. Doubling memory while staying inside the same regime changes only the run sizes and gains little.
- How can a query that requires ordered output avoid the sort operator altogether?By reading rows from a structure that already provides them in the required order - typically an ordered index whose key matches the requested ordering, possibly covering the needed columns. The plan then streams rows in order, eliminating the blocking sort, its memory grant and any temporary I/O, and it can also return the first rows immediately instead of waiting for the whole input.
Sorting a warehouse of paperwork with one small desk: sort a deskful at a time into labelled stacks, then walk the stacks taking whichever top sheet comes first. A bigger desk means fewer stacks to walk.
saying these in an interview costs you the question
- Describing spilling as random-access sorting on disk rather than sequential runs plus a merge
- Claiming a sort can emit its first row before consuming all input, ignoring that sorting is blocking (Top-N with a limit is the special case)
- Thinking temporary sort files are written to the write-ahead log or must survive a crash
- Assuming more sort memory always scales performance linearly instead of in steps
- Treating every spill as a memory problem when a bad row estimate or a missing ordered index is the real cause