Explain how a relational engine sorts a data set that does not fit in the memory available to the sort operator. Describe the algorithm and what determines its I/O cost.
answer
- Phase 1: fill, sort, dump run
- Phase 2: heap-merge F runs at a time
- runs ≈ N/M, passes = log_F(N/M)
- ≈ 2N I/O per pass — usually one pass
- Best sort = no sort (ordered index / partial sort)
basics
~20 sIt uses external merge sort: fill memory, sort that chunk, write it out as a sorted run, repeat; then merge the runs with a heap, taking the smallest current row across runs. Cost is roughly two I/Os per row per pass, and passes grow when runs exceed the merge fan-in.
solid answer
~60 sExternal merge sort has two phases. **Run generation.** Read input until the sort's memory budget is full, sort that chunk in memory (quicksort, or a replacement-selection heap), write it to a temporary file as a sorted *run*, and repeat until input is exhausted. With N pages of data and M pages of memory you get about N/M runs. **Merge.** Open several runs at once, keep one input buffer per run plus an output buffer, and repeatedly emit the smallest head row using a small heap. The number of runs merged at once is the *fan-in*, bounded by memory divided by buffer size. If the run count exceeds the fan-in, you merge in multiple passes, writing intermediate runs each time. Cost is dominated by passes: roughly `2 * N * (1 + ceil(log_F(N/M)))` page I/Os, where F is fan-in. In practice one or two passes covers almost everything, so the useful mental model is "data gets written and read about twice". More memory helps mainly by making runs longer and fan-in wider — and by avoiding the spill entirely.
code
text · 2 linesSORT (key=created_at, method=external merge, runs=14, temp_written=412 MB)
-> SEQ SCAN orders (rows=8,000,000, width=54)go deeper
Know the two phases by name — make sorted chunks (runs), then merge them — and that this is what "external" means.
Give the run-count and pass-count relationships, explain fan-in, and state the ~2N-per-pass I/O intuition.
Connect it to tuning: narrow the tuples before the sort, check whether an index supplies the order, and explain why extra memory has diminishing returns once you are at one merge pass.
Reason about it as a memory/IO budget allocation across a concurrent workload, and about which plan shapes preserve interesting orders so later operators avoid re-sorting at all.
## The problem Sorting is a *blocking* operator: it cannot emit its first output row until it has consumed all input, because any not-yet-seen row could belong first. If the whole input fits inside the sort operator's memory budget, the engine does an ordinary in-memory sort (typically quicksort, or an insertion sort for tiny inputs) and streams the result out. The interesting case is when it does not fit — the engine must sort correctly using bounded memory plus temporary disk space. The standard answer, essentially unchanged since the 1970s, is **external merge sort**. ## Phase 1 — run generation The operator reads input rows until its memory is full. It sorts that in-memory chunk and writes the sorted chunk to a temporary file. That sorted chunk is called a **run**. Then it clears memory and repeats. When the input ends, disk holds a set of sorted runs whose concatenation is the whole input. If the input is N pages and memory is M pages, this yields about `N/M` runs, each about M pages long. Some engines use **replacement selection** instead of sort-then-dump: maintain a heap of M pages of rows and continuously emit the smallest row that is still ≥ the last emitted one, refilling from input. On random input this produces runs averaging **2M** — half as many runs — and on already-nearly-sorted input it can produce one giant run. The trade-off is worse cache behaviour and more complex code, so several modern engines prefer simple quicksorted runs. A detail that matters in practice: what is written is not the base table rows but the **tuples the plan carries** — the sort key plus whatever columns downstream operators need, plus per-row overhead. Projecting away unused wide columns before the sort directly shrinks N. ## Phase 2 — merging Merging takes K sorted runs and produces one sorted stream. Memory is divided into one input buffer per run plus one output buffer, so the **fan-in** F is roughly `M / buffer_size - 1`. The operator reads the first block of each run, puts each run's current head row into a small heap (a *tournament tree* / loser tree of size K), and repeatedly pops the global minimum, refilling from whichever run it came from. Each output row costs one heap operation, O(log K) comparisons. If `K <= F`, one merge pass suffices and the merged output can be **pipelined** straight into the parent operator — no extra write. If `K > F`, the operator merges F runs at a time, writing longer intermediate runs, and repeats. The number of merge passes is `ceil(log_F(N/M))`. ## The cost formula Total page I/O is approximately: `2 * N * (1 + ceil(log_F(N/M)))` The leading `2 * N` is the run-generation write plus the first read; each further pass costs another read and write of the whole data. CPU cost is `O(N log N)` comparisons overall, and comparison cost is not free — sorting on a long text column with a locale-aware collation can cost more than the I/O. The practical shape of this: with even a modest budget and a fan-in in the tens or hundreds, `log_F(N/M)` is 1 for almost any realistic data size. So the honest rule of thumb is that a spilling sort writes and reads the data roughly **twice**, and only truly enormous sorts with tiny memory reach three or more passes. This also explains a counter-intuitive tuning fact: doubling the memory budget of a spilling sort rarely halves its time, because you were already at one merge pass. The big win is either avoiding the spill entirely or avoiding the sort entirely. ## Ways to avoid the sort - **Ordered access path.** Scanning an index whose key matches the required order supplies sorted rows directly; the sort operator disappears from the plan. This is the single biggest lever for ORDER BY on hot paths, at the cost of maintaining the index and, for non-covering indexes, random I/O to fetch the rest of the row. - **Order already established upstream.** A merge join or a sort-based aggregate leaves its output ordered on the join/group key; a later operator needing that same order can reuse it. - **Top-N.** If only the first few rows are wanted, the engine keeps a bounded heap instead of sorting everything. - **Partial sort / incremental sort.** If input is already sorted on a prefix of the required key, the engine sorts only within each group of equal prefix values, which bounds memory and starts producing rows early. ## Interview signals A strong answer names both phases, says what determines the number of runs (memory) and the number of passes (fan-in), gives the ~2N-per-pass cost intuition, and finishes with the observation that the fastest external sort is the one that never runs, because an index or an upstream operator already delivered the order.
- You double the sort operator's memory budget but a spilling sort barely gets faster. Why?Because it was almost certainly already doing a single merge pass. Doubling memory halves the run count and widens fan-in, but `ceil(log_F(N/M))` was 1 before and is still 1, so the total I/O stays around 2N. The step changes come from crossing back into a fully in-memory sort, or from removing the sort altogether with an ordered access path.
- How can an engine start returning rows before the whole sort finishes?Only in restricted forms. A final merge pass can be pipelined, so rows flow out as they are merged rather than being materialized again — but run generation must still have consumed all input first. An incremental/partial sort goes further: when the input is already ordered on a prefix of the sort key, the engine sorts each prefix group separately and emits it immediately, which bounds memory and gives a fast first row — very useful under LIMIT.
Merging library returns: you alphabetize one cartful at a time onto separate shelves (runs), then walk the shelves with one finger on each, always taking whichever book comes first alphabetically.
saying these in an interview costs you the question
- Claiming the engine "just sorts in memory and the OS pages it out" — swapping is not the mechanism, temp files are
- Saying the number of merge passes depends on data size alone, ignoring memory and fan-in
- Assuming each merge pass costs a random I/O per row; runs are read sequentially in blocks
- Believing a sort can emit its first row before consuming all input (outside of incremental/partial sort)
- Treating more sort memory as always the fix, when an ordered index removes the operator entirely