skip to content

How can a cluster produce an end-to-end ordered result without collecting every record onto one worker to sort?

level: middleimportance: must knowfreq 65%

answer

  1. make the placement carry the order
  2. you do not know the distribution yet
  3. boundaries from a sample, one span each
  4. sort locally, then read in piece order

basics

~20 s

Draw a sample of the ordering key, pick boundary values from it so each piece receives one contiguous span of the range, redistribute records by comparing the key against those boundaries, sort each piece in parallel, and read the pieces in index order.

solid answer

~50 s

The trick is to make the placement itself order-aware, so the final assembly is a read rather than a sort. First you need to know roughly how the keys are distributed, and you do not, so you sample the ordering key and take boundary values from the sample — say three boundaries for four pieces. Those boundaries become the routing function: instead of hashing the key, you compare it against the boundaries and send the record to the piece whose span contains it. Then every worker sorts its own piece, in parallel, and because piece `i` holds only keys below everything in piece `i+1`, reading the pieces in index order is the ordered result. No worker holds more than its own piece at any point. Runtimes differ in whether they derive the boundaries for you or expect you to supply them.

code

python · 25 lines
python
import bisect

# Boundary values chosen from a sample of the ordering key.
# 3 boundaries -> 4 pieces, each holding one contiguous span:
#   piece 0: key <= 120
#   piece 1: 120 < key <= 355
#   piece 2: 355 < key <= 810
#   piece 3: key > 810
boundaries = [120, 355, 810]


def piece_for_ordered(key):
    # left-biased: a key equal to a boundary belongs to the span ending there
    return bisect.bisect_left(boundaries, key)


def piece_for_hashed(key, pieces=4):
    # even spread, but the piece index says nothing about key order
    return hash(key) % pieces


assert piece_for_ordered(119) == 0
assert piece_for_ordered(355) == 1
assert piece_for_ordered(356) == 2
assert piece_for_ordered(9001) == 3

go deeper

for a junior

Recall the three moves in order: look at a sample of the keys, give each piece one contiguous span of the range, sort each piece on its own worker. The pieces are then read in order.

for a middle

Explain why the boundaries have to be derived from data rather than assumed, why they must be fixed before any record moves, and why the final assembly is a read rather than a sort.

for a senior

Judge whether the result is large enough to justify a sample plus a full movement at all, and say which parts of the packaging differ between runtimes rather than describing the one you have used.

for a principal

Weigh a recurring order-aware movement against publishing a weaker contract — ordered within a piece, or ordered by a coarse span — and count what the stronger contract costs on every run for every consumer that never needed it.

## The shape of the answer An end-to-end ordered result over more data than one machine can hold is produced by making the **placement** carry the order, so that sorting stays local and assembly is free. Concretely: 1. **Sample the ordering key.** The ordering key is the field the result must be ordered by. You need an approximate picture of how its values are distributed, and no such picture exists before the data is read, so one is built from a sample. 2. **Choose range boundaries.** Pick `n-1` key values from the sample as cut points, so that the key range is divided into `n` spans holding roughly equal numbers of records. These are **range boundaries**: the key values that separate one piece's span from the next. 3. **Redistribute using the boundaries as the routing function.** A **redistribution** is the step where every worker sends each record it holds to whichever worker will handle that record's key. Here the rule is a comparison against the boundaries, not a hash. 4. **Sort each piece locally, in parallel.** Every worker sorts only what it received. 5. **Read the pieces in index order.** Piece 0 holds the lowest span, piece 1 the next, and so on, so the concatenation is ordered. Nothing sorts at this step. ## Why a sample, and why not something exact Good boundaries are ones that split the data into equal shares, and that depends on the actual distribution of keys, which is data you have not read yet. The alternatives are worse: - **Fixed boundaries chosen by hand** work only if you already know the distribution and it is stable; they go stale the first time the data shifts. - **An exact scan to compute quantiles** costs a full pass over the input purely to decide placement. - **A sample** buys an approximation of the distribution for a fraction of that, and the approximation is good enough because the consequence of being slightly wrong is unevenly sized pieces, not a wrong answer. The sample has to be representative of the whole input, not of whichever part is cheapest to read — that is the failure mode this design is most exposed to. ## Range boundaries against a hash, side by side | | routing by a hash of the key | routing by range boundaries | |---|---|---| | equal keys meet | yes | yes | | piece index relates to key order | no | yes, by construction | | needs to know the distribution first | no | yes, hence the sample | | evenness of pieces | good, essentially by design | only as good as the sample | | final assembly of an ordered result | a sort somewhere | a read in piece order | ## What this costs - **The sample must precede the movement**, because the boundaries are an input to the routing function. Some designs pay a separate read of the input for it; others fold it into work that is happening anyway. The invariant is the ordering of events, not the price. - **The movement is a full redistribution**: every record crosses the network once, and the unit of that cost is bytes, not rows. - **The local sorts have a working set.** If a piece's sort exceeds what the worker can hold, part of the working set is written to local disk — a separate subject. - **The saving** is the single-worker pass you did not do: no machine ever holds the whole result, and every sort runs in parallel. ## What varies between runtimes, and what does not The mechanism above is the class-wide answer, but its packaging is not: - some runtimes derive the boundaries themselves whenever a program asks for an ordered result; - some expect the author to supply boundary values, or a comparator plus a sampling step, explicitly; - some express the sample as a visible extra step in the program graph, others hide it; - and over an **unbounded input** the whole idea narrows sharply: there is no complete result to order, so ordering is offered inside a bounded chunk — a grouping interval that has closed — or not at all. A candidate who says "the engine samples and does it for you" is describing one product. A candidate who says "the boundaries have to come from somewhere, and where they come from differs" is describing the class. ## The property that makes it work The entire construction rests on one invariant: **every key in piece `i` sorts before every key in piece `i+1`**. Hold that, and local sorts compose. Break it — by routing on a hash, by reusing boundaries over a differently distributed input, or by letting equal keys straddle a boundary — and the pieces stop composing, and you are back to a single-worker sort.

  • Why must the boundaries be fixed before any record is moved?
    Because the boundaries are the routing function. A record's destination is decided by comparing its key against them, so there is no destination to send it to until they exist. That ordering of events is what forces a sample ahead of the movement, and it is why the sample cannot be refined once records are in flight.
  • How many boundaries do you need, and what decides that?
    One fewer than the number of pieces: three cut points make four spans. How many pieces to use is a separate decision, driven by available parallelism and the size of each piece's working set, and it is owned by the question of how the input is cut up rather than by ordering itself.
  • Does routing by range boundaries still guarantee that equal keys meet?
    Yes. A span is defined by key comparisons, so two records with the same key compare the same way against every boundary and land on the same worker. That matters, because it means an order-aware routing also satisfies anything that needed grouping by key — you do not pay for both.

Sorting a mailroom's post by surname: rather than one clerk sorting the whole sack, you first glance through a handful of letters to see where the names actually fall, then label pigeonholes A to F, G to L and so on, drop each letter into its pigeonhole, let one clerk alphabetise each pigeonhole, and finally read the pigeonholes left to right. The labels must come from what the post actually contains — labelling by equal slices of the alphabet buries whoever gets the pigeonhole full of common surnames.

saying these in an interview costs you the question

  • Says the engine samples and builds the boundaries for you, as if every runtime did.
  • Picks boundaries by cutting the key range into equal-width spans regardless of the data.
  • Thinks a final merge on one worker is still needed after the pieces are sorted.
  • Believes the sample can be taken after the records have already been routed.
  • Assumes an end-to-end ordered result is available over an endless input.