Roadmap To Be A Data Engineer / Lesson 08
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)
- Expand — producer emits both:
item_conditionkeepsA/B/C, newcondition_gradecarriesexcellent/good/fair. Contract 1.4.0 → 1.5.0, additive, nothing can break. - 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.)
- 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
- 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?
- 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.
- Do we have a
CASE … ELSEin pricing, categorisation or reporting where theELSEis a silent default rather than a loudNULL? What’s the blast radius if every row lands there? - Which of our models
SELECT *from a source another team owns? - Last time we changed a column another team consumed — expand-and-contract, or just changed it?
- Do we monitor any shape metric — distinct values, branch mix, blended rate — or only volume, freshness and revenue?
8Hands on (~50 min)
- 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. - Write the contract. Save the YAML as
contracts/listings.yaml. Writevalidate(rows, contract)returning violations: missing required field, wrong type, value outsideenum, null rate over threshold. - Fail the build. Pre-release rows → clean. Post-release rows → one violation naming
item_conditionand the three offending values. That function in CI is the lesson. - 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. - Expand and contract. Emit both
item_conditionandcondition_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"))