skip to content

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%

answer

  1. never sort; match on the grain
  2. full outer join, four classes
  3. tolerance: absolute floor plus relative
  4. bucket aggregates first, descend second
  5. assert uniqueness before joining

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.

solid answer

~50 s

Never sort both sides: a total ordering forces every record to the worker owning its range and can cost more than the job under test. Instead declare the grain, full-outer-join the two results on it, and bucket every key into only-in-previous, only-in-latest, differs and unchanged, comparing numeric measures with a tolerance that has both an absolute floor and a relative part. At billions of rows, run it in two passes: first compare counts and sums per coarse bucket such as day or key prefix, then descend into only the buckets whose aggregates disagree. Report counts by class plus a bounded sample of differing keys, never a raw dump. Note the scope: this compares consecutive runs of the same logic. Comparing two versions of the logic before a cutover is a related but separate exercise, and there a non-empty diff is the expected result.

code

sql · 24 lines
sql
-- Both results carry one row per (day, account_id): the declared grain.
-- Rows are matched by that key, never by position, and measures are
-- compared with an absolute floor plus a relative allowance.
SELECT
    verdict,
    COUNT(*)                                 AS keys,
    SUM(latest_amount - previous_amount)     AS net_delta
FROM (
    SELECT
        CASE
            WHEN p.account_id IS NULL THEN 'only_in_latest'
            WHEN l.account_id IS NULL THEN 'only_in_previous'
            WHEN ABS(COALESCE(l.amount, 0) - COALESCE(p.amount, 0))
                 > 0.01 + 0.0001 * ABS(COALESCE(p.amount, 0)) THEN 'differs'
            ELSE 'unchanged'
        END                          AS verdict,
        COALESCE(p.amount, 0)        AS previous_amount,
        COALESCE(l.amount, 0)        AS latest_amount
    FROM previous_run p
    FULL OUTER JOIN latest_run l
           ON p.day = l.day
          AND p.account_id = l.account_id
) classified
GROUP BY verdict;

go deeper

for a junior

Recall that two large outputs are compared by joining them on a key that identifies a row, not by reading them in order, and that numbers need a tolerance.

for a middle

Explain the four result classes a full outer join produces, why a total ordering is the expensive way to get order-independence, and what a tolerance must contain.

for a senior

Show the two-pass shape — coarse bucket aggregates first, keyed comparison only where they disagree — and name what breaks it: duplicate grains, null keys and volatile columns.

for a principal

Make it a platform capability rather than a script per pipeline: one declared grain, one tolerance policy and one report shape, so a diff means the same thing across every dataset the organisation publishes.

## What this diff is, and what it is not This is a comparison of **two consecutive production runs of the same logic** — last night's published output against the run before it — asked because someone needs to know what changed and whether the change is explainable. It is not the same exercise as diffing **two versions of the transformation logic** over the same input before a cutover, which is a neighbouring subject and where a non-empty diff is the intended outcome rather than a problem. Same engine, same code, different run: that is the case here. ## Match on the grain, never on position 1. **Declare the grain** — the column set that identifies exactly one row. If both sides do not agree on it, the diff is meaningless before it starts. 2. **Full-outer-join** the two results on the grain, so unmatched keys on either side survive the join instead of vanishing. 3. **Classify every key** into one of four buckets, and report the counts: | class | what it usually means | |---|---| | only in previous | keys the latest run stopped producing — a lost input source, a tightened filter, or a genuine end of activity | | only in latest | new keys — new activity, a loosened filter, or a fan-out defect | | differs | the same key with a measure outside tolerance: the interesting class | | unchanged | within tolerance on every compared measure | 4. **Compare numbers with a tolerance** carrying an absolute floor and a relative allowance, so cents and billions are both handled by one rule. 5. **Emit a bounded sample** of differing keys with both values and the delta, plus the summary counts. A raw dump of a differing set is unreadable and, at this size, expensive to produce. ## Why not just sort both sides Sorting a large result into a total order forces a **redistribution** — the point where a step needs records currently held by other workers, so every worker writes its records out and every worker fetches the ones addressed to it — followed by a sort on each receiving worker. It can easily cost more than the job that produced the data, and it still leaves you comparing text or approximate numbers for exact equality. Matching on the grain gives order-independence for the price of one join, which is itself a redistribution but a single one, over the keys rather than over an imposed ordering. ## Aggregate first, descend second At billions of rows, run the diff in two passes: - **Pass one** computes, per coarse bucket — a day, a region, a hash prefix of the key — the row count and the sum of each compared measure on both sides. This is two grouped scans and produces a table small enough to read. - **Pass two** re-runs the keyed comparison **only for buckets whose pass-one aggregates disagree** beyond tolerance. Most nights, pass one is clean or disagrees in one bucket, and the expensive keyed join never runs at full width. A per-key canonical hash of the compared columns is a useful middle step when there are many measures: hash once per row, compare hashes, and fetch the actual values only for keys whose hashes differ. ## What breaks this diff - **A grain that is not unique.** Duplicates on either side turn the join into a cross product within the key. Assert uniqueness first, or compare counts per key as a multiset. - **Null keys.** A null join key matches nothing and lands in both only-in classes at once, which reads as a huge change. Exclude or substitute them explicitly, and count them separately. - **Columns minted during the run** — load timestamps, run identifiers, row numbers handed out during writing. Every row differs unless these are excluded. - **Comparing emissions instead of settled values.** Against a runtime that revises and re-emits a group, the comparison must read the latest value per key, not the stream of updates. - **A tolerance so wide it hides the class you care about.** The tolerance covers arithmetic reordering of approximate types, not business change. ## What the diff proves It proves the two runs agree or shows exactly where they do not. It does not say which one is right — that judgement needs the reconciliation against the input and the invariants behind it, plus an explanation for whatever drift remains.

  • How would you compare twenty measure columns without twenty comparisons per row?
    Compute a canonical hash per row over the normalised compared columns and diff the hashes, which reduces the per-key comparison to one equality. Then fetch the real values only for keys whose hashes differ. Normalisation matters: rounding and null handling must be identical on both sides or every hash differs.
  • The only-in-latest class holds ten million keys. What do you check first?
    Whether the grain is still unique on the latest side and whether the key was constructed differently — a changed hash, a trimmed string, a different case or a null substituted for a missing value all shift every key at once. A genuine ten-million-key expansion of activity would show in the input counts too.
  • Why report counts by class rather than the differing rows themselves?
    The shape of the change is the diagnosis: a handful of differing keys is usually source-side, while every key differing by a constant factor is systematic. Counts by class plus a bounded sample give that shape immediately, and producing a full differing set at this size is itself an expensive job.

saying these in an interview costs you the question

  • Sorts both outputs globally so the rows line up
  • Compares row by row on position and calls a reorder a regression
  • Joins on a grain that is not unique, producing a cross product per key
  • Leaves run timestamps in the compared column set
  • Dumps every differing row instead of counts plus a bounded sample
  • Confuses this with diffing two versions of the logic before a cutover