A compliance job must emit a 40 GB audit dump in key order. The service does Files.readAllLines, then Collections.sort. The heap is 4 GB. The process dies with OutOfMemoryError on the read, before a comparator runs. Nobody “chose a slow sort.” They ran an in-RAM procedure on a file that does not fit.
External sort writes sorted runs that fit in RAM, then k-way merges those runs with a heap of heads. Same merge primitive as merge sort, at disk scale: you never hold the whole input as one array. You hold one chunk, then k current heads.
This post is that two-phase procedure. Families and the catalog live on the Algorithms Roadmap. It is not a Hadoop tutorial, not a query-planner deep dive, and not a reason to paste a production sorter into a service that should call sort or ORDER BY.
The file does not fit, the chunk does
In RAM, merge sort splits an array you already allocated. Here the allocation is the failure. The layout is a file (or a stream you can rewind via temp files). The procedure is:
- Create runs. Read as many records as fit in a chunk. Sort the chunk in memory. Write it as a sorted run on disk. Repeat until the input is exhausted.
- K-way merge. Open the runs. Keep the next unread record from each run. Repeatedly emit the smallest of those heads, then refill that run’s head.
Step 1 is the in-RAM sort you already have — List.sort / Arrays.sort on the chunk, or the merge-sort sketch if you are teaching the merge. Step 2 is merge sort’s combine, generalized from 2-way to k-way. A binary heap of heads makes “smallest of k” a log k job instead of a linear scan of k, which matters when you have hundreds of runs.
You do not need the whole file in a List. You need chunkSize records plus k buffers for the merge. readAllLines on a file larger than the heap is not a sort strategy.
Phase 1: sorted runs
Pick a chunk that actually fits: object overhead, the sort’s extra array, and headroom for GC. A 4 GB heap does not mean a 4 GB ArrayList. Then flush.
input file (keys only, 12 records; RAM holds 4):
19, 7, 12, 3, 15, 8, 1, 14, 9, 4, 11, 6
chunk 0: 19, 7, 12, 3 → sort → run0: 3, 7, 12, 19
chunk 1: 15, 8, 1, 14 → sort → run1: 1, 8, 14, 15
chunk 2: 9, 4, 11, 6 → sort → run2: 4, 6, 9, 11
Three sorted files. Each is small enough to have been sorted in RAM. The original 12-key file never existed as one array.
If the last chunk is short, it is still a run. One run means the file already fit; you are done after phase 1 and there is nothing to merge. Zero records emit nothing.
Note: Sort each chunk with a stable in-RAM sort, and on a tie during k-way merge prefer the lower run id (earlier chunk). The whole file order among equal keys then matches the input. Drop either rule and “stable external sort” is a slogan.
A walked k-way merge
Three runs. The heap holds one head per run. Pop the minimum, write it, push the next line from that same run. Empty run: that slot disappears.
run0: 3, 7, 12, 19
run1: 1, 8, 14, 15
run2: 4, 6, 9, 11
heap of heads: 3(r0), 1(r1), 4(r2)
pop 1(r1) emit 1 push 8(r1) heap: 3, 8, 4
pop 3(r0) emit 3 push 7(r0) heap: 7, 8, 4
pop 4(r2) emit 4 push 6(r2) heap: 7, 8, 6
pop 6(r2) emit 6 push 9(r2) heap: 7, 8, 9
pop 7(r0) emit 7 push 12(r0) heap: 12, 8, 9
pop 8(r1) emit 8 push 14(r1) heap: 12, 14, 9
pop 9(r2) emit 9 push 11(r2) heap: 12, 14, 11
pop 11(r2) emit 11 run2 empty heap: 12, 14
pop 12(r0) emit 12 push 19(r0) heap: 19, 14
pop 14(r1) emit 14 push 15(r1) heap: 19, 15
pop 15(r1) emit 15 run1 empty heap: 19
pop 19(r0) emit 19 run0 empty heap: ∅
output: 1, 3, 4, 6, 7, 8, 9, 11, 12, 14, 15, 19
Each pop is one output record. Each push is one sequential read from a run. Sequential reads are the I/O you want; jumping around inside a 40 GB file is the I/O you do not.
If you have more runs than you can merge at once (too many open files, or k heads would blow RAM), merge batches of runs into larger runs and repeat. Same algorithm, extra pass. Production sort utilities already pick k and spill files. You do not need to invent a new pass policy for a blog sketch.
Java sketch
Phase 1: fill a chunk, sort it, write a run file. List.sort on the chunk is the in-RAM cousin — Timsort on objects, which is a merge of runs, not a reason to paste recursive merge sort into the flush.
static int writeRuns(BufferedReader in, Path dir, int chunkSize)
throws IOException {
int runId = 0;
List<String> chunk = new ArrayList<>(chunkSize);
String line;
while ((line = in.readLine()) != null) {
chunk.add(line);
if (chunk.size() == chunkSize) {
flushRun(chunk, dir, runId++);
}
}
if (!chunk.isEmpty()) {
flushRun(chunk, dir, runId++);
}
return runId;
}
static void flushRun(List<String> chunk, Path dir, int runId) throws IOException {
chunk.sort(null);
Files.write(dir.resolve("run-" + runId), chunk);
chunk.clear();
}
Phase 2: a PriorityQueue of heads. The heap is the layout (Heaps); this procedure is “always emit the smallest remaining head.”
record RunHead(int run, String line) {}
static void mergeRuns(List<BufferedReader> runs, BufferedWriter out)
throws IOException {
PriorityQueue<RunHead> heads = new PriorityQueue<>(
Comparator.comparing(RunHead::line)
.thenComparingInt(RunHead::run));
for (int i = 0; i < runs.size(); i++) {
String line = runs.get(i).readLine();
if (line != null) {
heads.add(new RunHead(i, line));
}
}
while (!heads.isEmpty()) {
RunHead h = heads.poll();
out.write(h.line());
out.newLine();
String next = runs.get(h.run()).readLine();
if (next != null) {
heads.add(new RunHead(h.run(), next));
}
}
}
That is the whole k-way step. Real records are not always a String line; a length-prefixed binary row is the same queue with a different compare. Do not build a TreeMap of entire runs. Do not load run files back into one List and call sort — that is the OOM you started with, plus extra copies.
Note: Close readers in finally / try-with-resources in code you ship. The sketch omits it to keep the merge visible. Buffer sizes matter more than a clever comparator: one page-sized read per run, sequential writes on the output.
Complexity
Let n be the number of records, M the records that fit in a chunk, R the number of runs (≈ n / M), k the merge width (k ≤ R).
| Phase 1 time | R in-RAM sorts: O(n log M) |
| Phase 2 time | O(n log k) heap operations if one merge pass |
| Extra RAM | One chunk, then k head buffers — not n |
| I/O | About two full scans per pass (read input/runs, write runs/output) |
If k = R, one merge pass finishes the file. If you must merge in several rounds, multiply the I/O by the number of rounds. The comparison bound stays O(n log n) in the usual case because log M + log R is log n. The bill that kills the compliance job is I/O and RAM, not the missing log n in a whiteboard recitation.
Stability is the same rule as in-RAM merge: left (here: earlier run) on a tie, plus stable chunk sorts.
When not to hand-roll this
Skip a custom external sorter when:
- The file fits.
Arrays.sort/List.sort/ merge sort in RAM. External sort on a 2 MB export is cargo-cult of the name. - A database already holds the rows.
ORDER BY(and the engine’s filesort / index) is the procedure. Exporting 40 GB to sort in your JVM is the wrong system boundary. - A sort utility already does the spills. GNU
sort(and the like) chunk, spill, and k-way merge.sort -Sis a memory budget, not a Java class you should rewrite on a Friday. - The job is a cluster shuffle. MapReduce / Spark sort-merge is this idea at machine scale. This post will not sketch YARN, Hadoop, or a shuffle service.
Do not hand-roll a production external sorter unless that is the job. A one-off compliance transform might be the job. A SaaS feature that “should just sort the audit log” is almost certainly ORDER BY or sort.
Also skip it when you do not need total order: a hash partition, a search index, or a streaming aggregate may answer the product question without ever emitting a fully sorted file.
JDK: a heap of heads, not an ExternalSort class
There is no java.util.ExternalSort. The library pieces are the ones you already call in RAM:
- Chunk sort:
List.sort/Arrays.sort - K-way heads:
PriorityQueuewith a comparator on the record (and run id on ties) - I/O:
BufferedReader/BufferedWriter, orFileChannelif you are measuring
Prefer sort or the database when they already sit on the box. Reach for the sketch when you are teaching the merge, or when the file and the memory budget are yours and no utility is allowed.
sort -S 1G -o audit.sorted audit.dump
Primitive Arrays.sort inside a chunk is fine when the chunk is int[] and you do not need stability on those keys. Object chunks that must keep input order among ties want the stable path, same as the merge-sort post.
Do not pull in a distributed-compute stack to k-way merge three runs on one machine. Do not implement a heap; PriorityQueue is the heap.
Cheat sheet
Job: total order on a file / stream larger than RAM
Phase 1: sort RAM-sized chunks → write sorted runs
Phase 2: k-way merge; heap of current heads; sequential reads
Tie-break: earlier run wins if the whole sort must be stable
Time: O(n log n) comparisons typical; I/O dominates
RAM: chunk + k buffers, not the whole file
Cousin: in-RAM merge sort (2-way merge of array halves)
JDK: List.sort on the chunk; PriorityQueue of heads
Do not: readAllLines; rewrite GNU sort; start a cluster for one box
Do:
- Size the chunk for objects, sort scratch, and GC — then spill runs.
- Merge with a heap of heads; read each run sequentially.
- Use
sortorORDER BYwhen they already solve the job.
Don’t:
- Load the file into a
Listand hope the heap holds. - Re-sort concatenated runs in RAM after phase 1.
- Hand-roll a production external sorter because the algorithm has a name.
- Treat this post as a MapReduce tutorial.
Wrap-up
External sort is merge sort’s combine step when the input is a file you cannot allocate. Sort what fits, write runs, k-way merge with a heap of heads. The extra “array” is now a directory of run files. Databases and sort utilities already do this; the JDK gives you List.sort and PriorityQueue, not a turnkey filesort. Hand-roll it when the file and the memory budget are the assignment. Otherwise call the tool that already spills.
That is the last sorting procedure in this wave that is about not fitting. Next optional step is BFS: level order, unweighted shortest path, and why the frontier is a queue.