skip to content

Producing an Ordered Result

A result ordered end to end rather than only inside each piece: boundaries drawn from a sample so each piece holds a contiguous span, sorted locally, and the pieces then read in order.

on this pageshow

questions

4

A job sorts each of its 200 pieces locally and writes them out; why is the combined result not ordered end to end?

level: juniorimportance: must knowfreq 70%

answer

  1. ordering is a property of placement
  2. a sort never moves records between pieces
  3. hashing spreads evenly, scatters the order
  4. contiguous spans, then read in index order

basics

~20 s

Sorting inside a piece orders only that piece. Unless the routing that filled the pieces gave each one a contiguous span of the key range, piece 3 can hold keys that belong before piece 2's, so reading the pieces in order interleaves ranges.

solid answer

~40 s

A local sort is an ordinary in-process sort over the records one worker happens to hold, so 200 of them give 200 independently ordered runs, not one ordered result. Which records a worker holds was decided earlier, by the routing function that sent each record to a destination. If that function hashes the key — the usual choice, because it spreads work evenly — the piece index tells you nothing about where those keys sit in the order. Equal keys still meet, since equal keys hash alike, but that is the only ordering you get. To read the pieces back in order you need the routing itself to be order-aware: boundaries chosen in key space so each piece owns one contiguous span, sorted locally, and then read in piece order.

go deeper

for a junior

Recall that a sort inside one piece can only order the records that worker holds, and that which records it holds was decided earlier by the routing function.

for a middle

Explain why hashing the key spreads work evenly but destroys any link between key value and piece index, and what has to change in the routing for pieces to concatenate into an ordered result.

for a senior

Show judgement about when an end-to-end order is worth its redistribution at all, and say plainly that it is a claim about a complete result, so over an endless input it applies to a bounded chunk or not at all.

for a principal

Frame ordering as a contract the platform publishes to consumers: ordered-within-piece is cheap and often sufficient, end-to-end ordering is a recurring cost on every run, and picking the wrong one leaks into every downstream reader.

## What a local sort actually guarantees A job's input is cut into **pieces** — one contiguous share of the input that a single worker processes on its own — and each piece is handled by one **unit of work**: one worker computing one piece once, so a retry replaces exactly that unit. When a step sorts, what runs inside that unit is an ordinary sort over the records that worker is holding at that moment. It has no view of any other worker's records and no way to acquire one. So after 200 local sorts you hold 200 independently ordered runs. That is genuinely useful — each run can be streamed, merged or searched as an ordered sequence — but it is not an **ordered result**, because an ordered result is a claim about the output read as a whole: every record in run 0 sorts before every record in run 1, and so on. ## The routing function decides this, not the sort Records are in a given piece because of the **routing function**: a rule applied to each record's key that names the one destination for it, so that equal keys always meet. Most routing is built to spread work evenly, which in practice means deriving the destination from a hash of the key. A hash deliberately destroys any relationship between a key's value and where it lands: - keys that are adjacent in the ordering land in unrelated pieces; - the piece index carries no information about which span of keys the piece holds; - equal keys do meet, because equal keys hash alike — that is the only ordering a hash-based routing function gives you; - sorting afterwards cannot repair this, because the sort never moves a record to a different piece. That last point is the whole answer. Sorting is a step that happens *after* placement, and placement is what an end-to-end order is a property of. ## The two ways to get an ordered result | approach | what moves | who sorts | where it breaks down | |---|---|---|---| | Collect the result onto one worker and sort it there | the entire result, to one machine | one worker, on its own | one machine's memory and local disk; no parallelism; the whole result crosses the network into one place | | Cut the key range at boundaries, then sort each piece | each record once, to the piece owning its span | every worker, in parallel | the boundaries are only as good as the sample they were drawn from | The first is fine when the result is small — a few thousand rows for a report, or a final aggregate. It stops being an option as soon as the ordered result is a large fraction of the input, which is the case an interviewer is asking about. ## The pipeline the second approach runs 1. Sample the **ordering key** — the field the result must be ordered by — across the input. 2. Choose boundary values from that sample so that each piece will receive one contiguous span of the key range, with roughly equal numbers of records per span. 3. Redistribute: a **redistribution** is the step where every worker sends each record it holds to whichever worker will handle that record's key, here using a routing function that compares the key against the boundaries instead of hashing it. 4. Sort within each piece, in parallel, so each piece is an ordered run covering a span nobody else covers. 5. Read the pieces back in index order. This is a read and a concatenation, not a sort: no worker ever has to hold more than its own piece. ## Finite results, endless inputs, and what varies An end-to-end ordered result is a statement about a *complete* result, so it is a finite-input idea. Over an unbounded input there is no complete result to order; runtimes over endless input generally offer ordering only inside a bounded chunk, such as one grouping interval that has closed, or not at all. Runtimes also differ in how much of the pipeline above they do on your behalf: some derive boundaries from a sample whenever an ordered result is requested, some expect the author to supply the boundary values, and the older designs express the whole thing as an explicit extra step in the program. Do not assume the one you learned is the model. ## What it costs, and what it is not - The sample has to be drawn before any record moves, because the boundaries are an input to the routing function. Some designs pay an extra read of the input for it; others fold it into work already happening. Either way the boundaries are fixed before the movement starts. - The movement itself is a full redistribution: every record crosses the network once, priced in bytes rather than rows. - The local sorts have a working set, and if it exceeds what the worker can hold, part of it is written to local disk — a separate subject with its own trade-offs. - Nothing about this makes the *contents* wrong. Locally sorted pieces are correct data with a weaker ordering contract than the reader assumed.

  • If the pieces are only ordered internally, is the output wrong, or just differently ordered?
    The records are all correct and none are lost — only the ordering contract is weaker than a reader might assume. Anyone consuming a single piece gets a valid ordered run; anyone reading the whole output in piece order gets interleaved key ranges. It is a contract problem, not a data problem, which is why it survives review so easily.
  • Why not just collect the result onto one worker and sort it there?
    Because that worker's memory, local disk and single-threaded sort become the ceiling, and the entire result crosses the network into one place. It is the right answer when the result is small — a report of a few thousand rows — and the wrong one as soon as the ordered result is a large fraction of the input.
  • Does a hash-based routing function guarantee anything useful about placement?
    Yes: equal keys hash alike, so all records sharing a key land on the same worker. That is enough for grouping, for folding by key, and for ordering records inside each key. It is not enough for an order across keys, which needs each piece to own a contiguous span of the range.

saying these in an interview costs you the question

  • Says sorting every piece in parallel is enough, because concatenation then gives an ordered result.
  • Believes a routing function that balances load well also preserves the order of keys.
  • Thinks the only way to get an ordered result is to pull everything onto one machine.
  • Assumes the sort step can move a record into a different piece to fix the order.
  • Treats locally ordered pieces as corrupt data rather than a weaker ordering contract.
open as a page

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%

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.

open as a page

Two requests resemble end-to-end ordering — the 100 largest values overall, and each key's records in time order; does either need range boundaries?

level: middleimportance: should knowfreq 45%

basics

~20 s

Neither does. The top few are found by keeping the best 100 inside each piece and merging those candidates, and ordering inside a key only needs equal keys to land together, which any routing on that key already gives. Range boundaries serve a complete result read in order.

open as a page

Range boundaries drawn from a sample leave one piece holding most of the records; what went wrong and what does it cost?

level: seniorimportance: should knowfreq 50%

basics

~20 s

The sample misrepresented the key distribution, so the cut points do not split the data evenly. The output is still correctly ordered; what suffers is the run, which is now bounded by the one oversized unit of work and its memory pressure.

open as a page