skip to content

Checking the Answer

Deciding a result is right when it is far too large to read: totals reconciled against the source, invariants and bounds instead of exact values, and diffs ignoring order.

on this pageshow

questions

5

Why does comparing a distributed job's output line by line against a stored expected result fail even when every value is right?

level: juniorimportance: must knowfreq 70%

answer

  1. arrangement is not value
  2. order is a race between workers
  3. file count follows width, not size
  4. match on the grain, not on position
  5. tolerance: absolute floor plus relative allowance

basics

~10 s

A distributed run promises the right records, not an order, a file layout or bit-identical arithmetic. A line-by-line comparison asserts all three at once, so it fails on arrangement while every value is correct.

solid answer

~40 s

The work is cut into pieces and each worker finishes when it finishes, so the order rows land in is a race rather than a property of the logic. The number of output objects follows how wide the last step ran, not how large the answer is, and sums of approximate numbers depend on the order they were added. A useful comparison therefore declares the grain, matches rows by that key rather than by position, applies an absolute-plus-relative tolerance to numeric measures, and either normalises or excludes columns minted during the run such as load timestamps and run identifiers. What is left failing is a difference in the values themselves, which is the only thing the test was ever trying to assert.

go deeper

for a junior

Recall that a distributed result is a set of records, not an ordered file: match rows by a key you can name, and never assert the order or the file count.

for a middle

Explain the mechanics — pieces finishing in arbitrary order, a redistribution's fetch race, output layout following the width of the last step, and approximate addition not being associative.

for a senior

Show judgment about what the comparison proves: a blessed expected file catches change, not correctness, and needs a reconciliation or an invariant behind it before anyone relies on the number.

for a principal

Frame the standing rule: every published dataset declares a grain and a tolerance, so comparisons across the organisation mean the same thing and nobody re-invents a bespoke diff per pipeline.

## What a stored expected result actually asserts A **stored expected result** is a file of rows committed next to the test, and comparing the job's output to it line by line asserts three things at once: 1. the same **values** were produced; 2. they appear in the same **order**; 3. they are laid out in the same **shape** — the same number of objects, the same columns, the same text rendering of each number. Only the first is a statement about the job's logic. The other two describe how one particular distributed execution happened to arrange itself, and a distributed run promises neither. ## What a distributed run does not promise - **Order across workers.** The input is cut into **pieces** — one chunk of the data that one worker processes independently of the others — and each worker finishes when it finishes. Whichever worker's output is committed first sets the order in the combined result; next time a different one wins. A job that ran over a single piece during development can look perfectly stable and then reorder the first time it runs wide, which is why this bites in production rather than in the test. - **Order within a group.** A step that needs records currently held by other workers forces a **redistribution**: every worker writes its records out, and every worker fetches the ones addressed to it. The order those fetched records arrive in is a race between senders. - **The layout of the output.** Many runtimes write one output object per piece, so the file count tracks the width of the last step rather than the size of the answer; others commit one set through a single writer. Change the width and the identical answer arrives as a different set of files. - **Bit-identical arithmetic.** Approximate numeric types are not associative under addition: the same amounts added in a different grouping give a total that differs in the last places. The value is right to the precision the type carries, not to the last bit. *Why* two runs can disagree at all is a separate subject; here it is enough that they may. - **Columns minted during the run.** Load timestamps, run identifiers, and surrogate numbers handed out as rows are written differ on every execution by design. ## What "the output" even is depends on the model | the job | what the output is | what an expected-result comparison can mean | |---|---|---| | a finite run over bounded input | one complete set of records, committed at the end | compare the whole set, once | | continuous work run as a rapid succession of small finite runs | a series of increments, one per small run | compare one increment, or compare a time bucket once its completeness claim has passed | | a record-at-a-time runtime holding a key-bound retained set | a running result per key that may be re-emitted as the group is revised | compare the **latest value per key**, never the sequence of emissions | A test that captures an expected file from the first model and re-uses the comparison style against the third will fail forever, because the third legitimately emits the same key more than once. ## Comparing so that only the values matter 1. **Declare the grain** — the column set that identifies exactly one output row. If you cannot name it, you cannot diff; assert uniqueness on it as its own check. 2. **Match by key, not by position.** Join the produced rows to the expected rows on the grain and classify each key as present-in-both, only-produced, or only-expected. 3. **Compare numbers with a tolerance** that has both an absolute floor and a relative allowance, so small values are not swamped by the floor and large ones are not failed by the last digit. 4. **Normalise the volatile columns** — exclude them, or overwrite them with a fixed value before comparing. 5. **Decide the multiset question.** If the grain is genuinely not unique, compare counts per key rather than pretending each key has one row. ## What this style of comparison still cannot prove - A stored expected result is usually produced by running the job once and blessing the output, which freezes whatever the job did that day, bug included. It detects **change**, not **correctness** — which is why totals reconciled against the input and invariants asserted on the result carry the correctness argument instead. - It says nothing about volume. A comparison over a handful of rows proves the shape of the logic and nothing about what happens at production width. - Executing the whole job inside the test's own process, on one machine with no network and no second worker, cannot surface record redistribution, worker loss or recovery at all — that limitation, and how to work around it, is the neighbouring subject of testing logic outside the cluster.

  • Why not simply sort both sides before comparing them?
    A total ordering over a large result forces every record to move to the worker that owns its range — a full redistribution plus a sort, which can cost more than the job under test. It also does not help with tolerance, volatile columns or duplicate keys. Matching on the declared grain gives the same order-independence without the movement.
  • The grain has duplicates in both outputs. What changes?
    Compare as a multiset: aggregate each side to a count per key and diff the counts, or add a deterministic tie-break column that is part of the data rather than minted during the run. Silently deduplicating hides exactly the defect — an accidental fan-out — that the comparison should catch.
  • Does a single-piece run in a test make the order stable enough to assert?
    Inside that one execution, often yes — which is the trap. The assertion then encodes an accident of the test's width, and breaks the first time the same logic runs wide or a step redistributes. Write the comparison order-insensitively even when the current run happens to be ordered.

saying these in an interview costs you the question

  • Treats row order in a distributed output as stable enough to assert
  • Sorts both sides globally just to make a comparison work, at full redistribution cost
  • Blesses whatever the job printed once and calls the file an expected result
  • Compares approximate numbers for exact equality and calls the mismatch a bug
  • Leaves run timestamps and generated identifiers in the compared columns
  • Assumes one output object means the job produced one piece of work
open as a page

Which totals do you reconcile against the input to judge a nightly aggregation job, and what does a match still miss?

level: middleimportance: must knowfreq 62%

basics

~20 s

Carry a small set of totals through the job: rows in against rows out with the relation the transform implies, the grand total of an additive measure, the distinct count of the grouping key, and the null and rejected counts. A match proves the job faithful to its input, not the number true.

open as a page

When no independent expected value exists for an aggregate, which invariants and bounds still catch a wrong result, and which errors survive?

level: middleimportance: should knowfreq 55%

basics

~20 s

Assert what must be true of any correct result rather than the value itself: uniqueness at the declared grain, parts summing to the whole, keys present in the reference set, measures inside their physical range, and volume inside a band drawn from history. A plausible wrong number survives all of them.

open as a page

How do you decide whether a 0.4% move between last night's published total and the previous run's is real or a defect?

level: seniorimportance: should knowfreq 46%

basics

~20 s

Decompose the delta before judging it: attribute the move per key and per period, then check whether the input moved too. A few keys shifting is usually the source; every key shifting proportionally is usually the job. Explain the remainder rather than widening the tolerance.

open as a page

How do you diff last night's output against the previous run's when both hold billions of unordered rows at the same grain?

level: seniorimportance: should knowfreq 50%

basics

~20 s

Match on the declared grain rather than on position, classify each key as only-previous, only-latest, differing or unchanged, and compare measures with an absolute floor plus a relative allowance. Aggregate per bucket first and descend only into the buckets that disagree.

open as a page