Roadmap To Be A Data Engineer / Lesson 08

Lesson 08 Data contracts Fundamentals §1 Source Systems About 10 min read

The Clause Nobody Wrote

A source changes meaning without changing shape. Schema drift and contracts.

Sequel to Lesson 06 (a feed that died silently) and Lesson 07 (a DAG that went green on empty data). Both of those were failures inside our platform. This one starts outside it, in a table another team owns and is entirely entitled to change.

1The scene

Wednesday 19 August, 14:07 CEST. The backend team ships PR #4471 — grade taxonomy cleanup against shop.listings, the table our warehouse copies hourly. Four column changes; two matter.

# Change Kind
1 brand_id → brand_key (int → string) rename — structural
2 item_condition values A/B/C → excellent/good/fair re-coding — semantic
3 + condition_scored_at timestamptz additive
4 seller_note varchar(255) → text widening

Change 1 broke six dashboards at 14:11 — four minutes after deploy. Slack 14:19, patch merged 14:52. Nothing was ever wrong; things were absent. Cost: an annoying afternoon, €0.

Change 2 broke nothing. fct_price.sql, written eleven months earlier:

CASE item_condition
  WHEN 'A' THEN 1.00   -- as new: full price
  WHEN 'B' THEN 0.90   -- light wear: -10%
  WHEN 'C' THEN 0.80   -- visible wear: -20%
  ELSE 0.80            <-- the branch that ate the money
END AS markdown_multiplier

From 14:07 the column held 'excellent', 'good', 'fair'. No WHEN matched. Every one of 8,000 listings fell to ELSE 0.80 — a 20% markdown on everything, including the 2,080 grade-A pieces that should have sold at full price. Nine days, 2,285 orders, €29,289.95 under-charged (10.91% of intended GMV, €12.82/order). Every job green throughout.

A rename is a broken bone: it hurts immediately and you set it the same day. A re-coding is a slow bleed: nothing hurts, and you find it from the blood on the floor.

2Schema drift, classified

Drift sorts on two axes; only one is expensive.

Pipeline keeps running Pipeline breaks
Numbers change Silent & expensive — enum re-coded, unit changed (cents→euros), meaning changed, timezone changed. Name and type identical; nothing structural to detect. Loud & partial — type narrowed, new upstream NOT NULL. Some rows fail, some pass.
Numbers safe Free — column added, type widened, comment edited. The majority of upstream changes; absorb them silently. Loud & cheap — column renamed, dropped, retyped. The failure is the alert; no wrong number reaches anyone.

The instinct after an outage is to harden against the right-hand column, because that’s what woke people up. The money is in the left-hand column. An outage is a change that announced itself.

3Was it invisible? (the honest check)

Two signals existed for nine days. The series habit is to compute before claiming invisibility — here the answer is half:

  • Revenue: genuinely ambiguous. Median daily GMV fell 15.9% (z = −2.18), but 3 of the 9 post-release days land inside the pre-release min–max band [€27,186 – €36,784]. A 16% dip in late August is an argument, not an alarm.
  • Shape: unmissable. Markdown as a share of list price sat inside 9.99%–10.31% for 18 days, then hit exactly 20.00% and stayed there for nine days, to the second decimal. Blended rates over a real mix of goods do not land on a round number and hold.

Two more one-line checks, either of which ends it on day one:

Query 1–18 Aug 19–27 Aug
COUNT(DISTINCT markdown_multiplier) 3 1
share of orders on the ELSE branch 28% 100%
listings with item_condition outside ('A','B','C') 0 8,000
daily GMV vs the 18-day band — 3 of 9 inside

Revenue is an aggregate, and aggregates have weather. The shape of the data — distinct values, branch mix, blended rates — has no weather.

4The fix: a data contract

A machine-checkable agreement at the boundary between producer and consumers. Not a wiki page — a file in the producer’s repo, checked in CI, versioned like code.

# contracts/shop.listings.yaml
version:   1.4.0
owner:     team-marketplace-backend
consumers: [data-platform, pricing-service, finance-reporting]
freshness: {max_lag: 1h}          # lesson 06's threshold, written down
volume:    {min_rows: 6500, max_rows: 9500}

fields:
  - name: item_condition
    type: string
    enum: [A, B, C]               # <-- the clause nobody wrote
    meaning: >
      Grader's condition band, assigned at intake.
      A = as new · B = light wear · C = visible wear.
    required: true

  - name: price_cents
    type: integer
    unit: EUR_cents               # units are a promise, not a comment

evolution:
  additive: allowed               # new fields, widened types: minor bump
  breaking: expand_and_contract   # renames, re-codings: never in one release
  notice:   14d

Six clauses, each buying something:

Clause Answers Would it have caught PR #4471?
Owner & consumers who to talk to before changing this yes — review pings data
Schema which columns and types exist the rename, yes
Semantics units, allowed values, meaning the re-coding — this is the one
Guarantees nullability, uniqueness, freshness, volume no, but catches Lessons 06 & 07
Evolution rules what may change without a major version yes — bans the one-release rename
Enforcement where the check runs turns all of the above into a red build

With the enum clause, PR #4471 never merges: the contract test sees 'excellent' outside the declared domain and fails the build, naming the column, the offending values and the three consuming teams. Cost of the incident: one conversation.

5Expand and contract (how a breaking change is supposed to happen)

  1. Expand — producer emits both: item_condition keeps A/B/C, new condition_grade carries excellent/good/fair. Contract 1.4.0 → 1.5.0, additive, nothing can break.
  2. Migrate — consumers move on their own schedule inside the notice window. The contract’s consumer list plus lineage tells you who is left. (Impossible without knowing who reads you.)
  3. Contract — when nobody reads the old field, drop it as a major bump 1.5.0 → 2.0.0.

Same end state as PR #4471, reached with nobody’s prices wrong for nine days. Cost: one extra column for a fortnight. For event streams the same discipline is packaged as a schema registry (Avro/Protobuf + Confluent Schema Registry), enforcing BACKWARD/FORWARD/FULL compatibility before a producer may publish.

6Where to put the check — each rung costs ~10× the one above

Rung Check Cost when it catches
1 Producer CI — contract test on the migration/payload one conversation
2 Ingest gate — reject or quarantine violating rows one alert
3 Warehouse tests — accepted_values, not_null, unique, relationships one failed run
4 Shape monitors — distinct counts, branch mix, blended rates one day of bad prices
5 A human notices the money €29,289.95

And one habit worth more than any tool: never SELECT * across a team boundary. A wildcard forwards whatever arrives — which is exactly how a structural change becomes a silent one. An explicit column list turns a rename into a rung-1 build failure, for free.

7Ask your team

  1. For our three most important upstream tables — who owns them, and do they know we read them? If a backend engineer renamed a column tomorrow, would anything stop the merge?
  2. Where are the enum domains and units written down? Not the types — the allowed values, and whether a column is cents or euros, local time or UTC.
  3. Do we have a CASE … ELSE in pricing, categorisation or reporting where the ELSE is a silent default rather than a loud NULL? What’s the blast radius if every row lands there?
  4. Which of our models SELECT * from a source another team owns?
  5. Last time we changed a column another team consumed — expand-and-contract, or just changed it?
  6. Do we monitor any shape metric — distinct values, branch mix, blended rate — or only volume, freshness and revenue?

8Hands on (~50 min)

  1. Reproduce. Run seed.py (below). Confirm 8,000 listings, 2,080/3,680/2,240 grade mix, UNDER-CHARGED EUR 29,289.95. No RNG — your figures match to the cent.
  2. Write the contract. Save the YAML as contracts/listings.yaml. Write validate(rows, contract) returning violations: missing required field, wrong type, value outside enum, null rate over threshold.
  3. Fail the build. Pre-release rows → clean. Post-release rows → one violation naming item_condition and the three offending values. That function in CI is the lesson.
  4. Add the two shape monitors. Emit COUNT(DISTINCT markdown_multiplier) and the blended markdown rate daily; alert when the distinct count drops or the rate moves >1pp off its 14-day median. Show both fire on 19 Aug while a GMV-band alert stays quiet on 3 of 9 days.
  5. Expand and contract. Emit both item_condition and condition_grade; show an old and a new consumer both passing against 1.5.0, then drop the old field as 2.0.0.

Push as de-practice/08-data-contracts. README: name one real upstream table we depend on, its consumers, and the one enum or unit clause that is true today but written down nowhere.

9Takeaway

Structure breaks. Semantics lie. A pipeline detects the first on its own — a missing column is an error. It can never detect the second, because 'excellent' is a perfectly valid string. The only thing that catches a meaning change is a written-down promise about what the values may be, checked by a machine, in the producer’s build, before the change ships.

Smallest useful action: pick the one upstream column whose values your business logic branches on, and add an accepted_values test. One clause, five minutes, and the expensive half of schema drift becomes the cheap half.

10Vocabulary

  • Schema drift — any unannounced change to a source you don’t own: added, renamed, dropped, retyped or (worst) re-coded.
  • Structural vs semantic change — names and types vs units, allowed values and meaning. Structural breaks loudly; semantic doesn’t break at all.
  • Data contract — versioned, machine-checkable agreement between a producing team and its consumers: schema, semantics, guarantees, evolution rules, owner.
  • Enum domain / accepted values — the closed set a column may contain. The single clause that would have prevented this incident.
  • Expand and contract — three-step migration for breaking changes: emit both, let consumers move, remove the old one in a major version.
  • Schema registry — the streaming equivalent; rejects producer schema versions that break declared compatibility.
  • Shape monitoring — tracking distinct values, branch mix and blended rates over time; the signals that have no weather.
  • Quarantine table — where contract-violating rows go, so they are neither silently accepted nor silently dropped.

11Appendix — seed.py

Show the full code (72 lines)
"""
Lesson 08 - Schema drift & data contracts.  Deterministic: no random, no seed,
modular arithmetic only.  Every figure in the lesson is produced by this file.
"""
from datetime import date, timedelta
import json, statistics as st

N_ITEMS = 8000
DAYS    = 27                      # 2026-08-01 (Sat) .. 2026-08-27 (Thu)
DAY0    = date(2026, 8, 1)
RELEASE = 18                      # index of 2026-08-19: PR #4471 ships 14:07
MULT      = {"A": 1.00, "B": 0.90, "C": 0.80}   # the CASE branches in fct_price.sql
ELSE_MULT = 0.80                                # ... ELSE 0.80 END
DOW = {0:.88, 1:.92, 2:.95, 3:1.00, 4:1.05, 5:1.22, 6:1.18}   # Mon..Sun demand

def item(i):
    """price in cents, condition grade - both pure functions of the item id."""
    return 1290 + (i * 4703) % 23711, "ABC"[(0 if (i*17) % 50 < 13 else
                                             1 if (i*17) % 50 < 36 else 2)]

items  = {i: item(i) for i in range(1, N_ITEMS + 1)}
grades = {g: sum(1 for _, gg in items.values() if gg == g) for g in "ABC"}

daily, cursor = [], 0
for d in range(DAYS):
    day    = DAY0 + timedelta(days=d)
    wobble = 1 + ((d * 37) % 17 - 8) / 100                 # +/- 8%, deterministic
    n      = round(250 * DOW[day.weekday()] * wobble)
    rows   = []
    for _ in range(n):
        cursor += 1
        rows.append(items[(cursor * 2213) % N_ITEMS + 1])
    listed  = sum(p for p, _ in rows)
    correct = sum(round(p * MULT[g]) for p, g in rows)     # what should be charged
    drifted = sum(round(p * ELSE_MULT) for p, g in rows)   # every row hits ELSE
    live    = drifted if d >= RELEASE else correct
    daily.append({"date": day.isoformat(), "dow": day.strftime("%a"), "orders": n,
                  "live": live, "correct": correct, "lost": correct - live,
                  "else_share": 1.0 if d >= RELEASE else
                                sum(1 for _, g in rows if g == "C") / n,
                  "n_mult": 1 if d >= RELEASE else 3,
                  "markdown_pct": (1 - live / listed) * 100})

pre, post = daily[:RELEASE], daily[RELEASE:]
gmv  = lambda xs: [x["live"] / 100 for x in xs]
lost = sum(x["lost"] for x in post)
want = sum(x["correct"] for x in post)
n    = sum(x["orders"] for x in post)
band = (min(gmv(pre)), max(gmv(pre)))
md   = lambda xs: [x["markdown_pct"] for x in xs]

print(f"items {N_ITEMS}  grades A={grades['A']} B={grades['B']} C={grades['C']}")
print(f"window {post[0]['date']}..{post[-1]['date']}  {len(post)} days  {n} orders")
print(f"UNDER-CHARGED  EUR {lost/100:,.2f}   {lost/want:.2%} of intended   "
      f"EUR {lost/n/100:.2f}/order")
print(f"signal 1 revenue: pre band [EUR {band[0]:,.0f} .. EUR {band[1]:,.0f}], "
      f"post median {st.median(gmv(post))/st.median(gmv(pre))-1:+.1%}, "
      f"z {(st.median(gmv(post))-st.mean(gmv(pre)))/st.stdev(gmv(pre)):+.2f}, "
      f"{sum(1 for v in gmv(post) if band[0] <= v <= band[1])}/{len(post)} days inside band")
print(f"signal 2 shape:   markdown pre [{min(md(pre)):.2f}% .. {max(md(pre)):.2f}%] "
      f"-> post {post[0]['markdown_pct']:.2f}% flat;  "
      f"distinct multipliers {pre[0]['n_mult']} -> {post[0]['n_mult']};  "
      f"ELSE branch {st.median(x['else_share'] for x in pre):.0%} -> "
      f"{post[0]['else_share']:.0%}")
print()
for x in daily:
    print(f"{x['date']} {x['dow']} n={x['orders']:3d} "
          f"charged=EUR{x['live']/100:9,.2f} lost=EUR{x['lost']/100:8,.2f} "
          f"markdown={x['markdown_pct']:5.2f}%"
          + ("   <-- PR #4471" if x["date"] == post[0]["date"] else ""))

json.dump(daily, open("daily.json", "w"))
Back to top