skip to content

Two runs of one job over the same input write the same rows in a different order — why?

level: juniorimportance: must knowfreq 65%

answer

  1. order is not part of the contract
  2. many pieces, whichever finishes first
  3. a redistribution discards arrival order
  4. only a declared sort promises row order
  5. first-per-key needs a tiebreaker field

basics

~20 s

Nothing in the job asked for an order, so none was kept. The work is cut into pieces processed by many workers at once, and rows appear in whatever sequence those pieces finished and were collected. Only a declared ordering step makes row order stable.

solid answer

~50 s

A job's handle on its data is a distributed collection — the engine's view of records spread across the cluster — and a set of records has no sequence. Output order is a by-product of three things the program never stated: which piece each record landed in, which pieces finished first, and the order a collecting worker received them. So the honest default claim is *the same rows, ignoring order*, not *the same rows in the same order*. A step that needs records held by other workers redistributes them (a shuffle), which discards whatever order the input had; where an engine sorts grouping keys as part of that redistribution the output can look ordered, but that is a property of how that engine redistributes, not a promise of the model. If you need an order, declare one — and compare two runs as multisets, or sort both first.

go deeper

for a junior

Recall that a job promises which rows it produces, not the sequence they come out in. If an order matters to the consumer, the job has to declare one.

for a middle

Explain where the order comes from: which piece a record landed in, which pieces finished first, and what a redistribution between workers did to arrival order.

for a senior

Show that you compare runs as multisets or after an explicit sort, and treat first-record-per-key without a tiebreaker as a correctness bug rather than a presentation detail.

for a principal

Frame it as a platform contract: which outputs may claim a stable order, what producing that order costs, and where consumers must be stopped from depending on an incidental one.

## Say which equality you mean "The two runs produced the same result" hides at least four different claims, and a job promises different things about each. | claim | what it asserts | promised by default | |---|---|---| | same rows, ignoring order | the outputs are the same multiset of records | usually yes, for a job with no unstable source inside a step body | | same rows in the same order | row *n* of one output equals row *n* of the other | **no**, unless the job declares a total order | | bit-identical values | every number matches to the last digit | **no**, where the order values were combined in affects the arithmetic | | same output files, same names | the output is laid out identically in the shared store — the cluster-visible storage a job writes to | **no**; how many output pieces there are follows the run | Row order is the second line, and it is the one people assume they already have because a small test on one machine appeared to give it to them. ## Where output order actually comes from The engine derives a **step graph** from the program — the ordered set of steps, each naming the steps whose output it reads — and runs each step over the input in pieces, one worker thread to a piece. Nothing in that description mentions a sequence. The order rows land in the output is set by: - **which piece each record was in** — fixed by how the input was cut, and then changed by any step that needs records held by other workers, which redistributes records so that all records sharing a grouping key meet on one worker (a **shuffle**); - **which pieces finished first** — a worker that finishes early contributes early, and a piece that runs far longer than its peers (a **straggler**) contributes last; - **the order the collecting worker, or the coordinating process that receives whatever the demanding call asked for, took delivery.** None of those three is stated anywhere in the program, so none of them is promised. ## The order that looks stable until it isn't Many jobs appear to produce a stable order for months. That is an **incidental** order, and it holds only while nothing moves. It changes when: 1. the step width changes — how many pieces a step processes at once; 2. a piece is lost and recomputed, so it arrives later than it did before; 3. the set of input files changes, or the same records are laid out differently; 4. the engine adjusts the not-yet-run part of the plan using measurements from work already finished. A consumer that quietly depends on the incidental order breaks on the day one of those four happens, which is rarely the day anyone is looking. ## Where engines genuinely differ - **Within one piece**, most engines feed records to a per-record function surface — a surface where the author hands over a function to call once per record — in the order the piece was read. That is a property of the piece, not of the output. - **Across a redistribution**, engines differ: some assign keys to destinations by hash and emit groups in no meaningful order; others sort keys as part of the redistribution, so their output looks sorted. Relying on the second is relying on one engine's method. - **In a continuous record-at-a-time model**, where one fixed graph stays running and each record passes through as it arrives, output order also reflects the interleaving of concurrent sources, which nothing reproduces. - **In a repeated-small-batch model**, which runs continuous work as a fast succession of small finite jobs, each small job emits its own block of output and the boundary between blocks moves with timing. ## Order inside a group is not promised either This is where unpromised order stops being cosmetic. A step that takes *the first record per key*, or *the last status per account*, is reading an order nobody declared: whichever record reached the grouping first this time wins. The rows are the same both runs; the **values** are not. The repair is to define "first" in the data — order by a timestamp with a unique field as a tiebreaker, then take the top — so the answer is a function of the records rather than of the race. ## Comparing two runs Compare as multisets, or sort both outputs by a key you choose before comparing. A line-by-line difference of two unsorted outputs reports thousands of differences that mean nothing and hides the one that means something. Checking an answer far too large to read at all is a separate subject with its own owner, and producing a genuinely ordered result at scale is another; here the point is narrower — state the equality you are claiming before you claim it.

  • Same code, same input, same number of pieces — why can the order still differ?
    Completion times still differ, so pieces reach the collector in a different sequence. A recomputed piece arrives later than it did before, and on engines that adjust the not-yet-run part of the plan from measurements of finished work, even the shape of the run can differ. Equal width makes the order more often incidentally stable; it never makes it promised.
  • A downstream step takes the first record per key. Is that affected by this?
    Yes, and more seriously than row order, because it changes values rather than presentation. Whichever record reached the grouping first wins, and that can differ between runs. Define the order in the data instead: sort by an event timestamp with a unique field as tiebreaker and take the top, so "first" is a property of the records.
  • Is the number of output files part of what a job promises?
    No. How many pieces the output is written as follows the run — the width of the final step and any rewrite the engine applied. Treat the output as a set of records at a location, not as a file whose name or count is stable, and never make a consumer read one specific file.

saying these in an interview costs you the question

  • Says output rows come back in input-file order.
  • Assumes a grouping returns its keys in sorted order.
  • Believes a fixed piece count guarantees a fixed output order.
  • Calls unordered output a defect in the engine.
  • Diffs two unsorted outputs line by line and reports failure.
  • Takes the first record per key with no tiebreaker.