A kitchen buys flour from a supplier. Each delivery comes with a note: a line that says "flour", then a number. The kitchen's system reads the note, checks that each line has a name and a number, and orders the next batch from it.
One day the supplier changes its scales. It now counts flour in pounds, not kilograms. The note still says "flour, 25". The name is right, the number is a number, and every check the kitchen runs passes. The kitchen now has less than half the flour it thinks it has, and nobody finds out until the bread comes out wrong.

Nothing failed in that story. The note had the right shape, so every check said yes. The problem was in what the number meant, and no check looked at meaning.
Data that feeds a model has exactly this problem. Another team writes the data, and your model reads it. They can change it for good reasons, on their own schedule, and your checks may only look at its shape. In this lesson I measure which changes a shape check stops and which ones walk straight into a model, on twelve realistic changes and real data.
This is the tenth lesson of the chapter on data engineering for machine learning. Three earlier lessons touch the same ground, and I build on them rather than repeat them.
Lesson 2, on data ingestion pipelines, renamed a column upstream and watched pandas fill it with empty values that became 0, with no error. A check on column names caught it. Lesson 5, on data validation, measured range, category and consistency checks on real weather data, and their false alarms. Lesson 6, on data versioning, kept versions of the data itself.
This lesson is about the step before all of that. It is the agreement between the team that writes the data and the teams that read it. It is also about the tools that check a change to that agreement before it ships.
What this rebuild corrected. An earlier version of this lesson made claims its sources do not support. Each is fixed where it comes up, and the full record is in
results/ctr-factcheck.json.
- It said a schema registry "rejects events that break the schema". A registry checks schemas, not messages.
- It said a rename is safe in no compatibility mode. Avro readers can rename with an alias, a second name the reader accepts for a field. Protobuf's binary format does not carry names. And the lab below shows a rename passing every mode, which is worse, not safer (see the slide on defaults).
- It said BACKWARD is "the safest" mode. It is the default; FULL_TRANSITIVE is the strictest.
- It said "most large organizations" use FULL, and that Airbnb's work was about unannounced upstream changes. I found no source for either, so I cut the first and corrected the second.

The plan runs from top to bottom. First what a contract is, and the rules a registry checks. Then the lab, which replays twelve changes and measures what each mode let through. Last, what a shape check cannot see, and how to make a change that no mode allows.

The producer is the team, and the code, that writes the data. A consumer is anyone who reads it, such as the job that builds a model's inputs. A record is one row of data, for example one hour of bike rides. A field is one named value inside a record.
A schema is the written shape of a record: the names of its fields and the type of each, such as whole number or text. A data contract is wider. It is the producer's written promise about the data: its shape, what each value means, how fresh it will be, and who owns it.
A schema registry is a service that keeps every version of a schema and checks each new one before it is accepted. Compatibility says whether a reader on one version can read data written with another. BACKWARD, FORWARD and FULL are the three main kinds, and the next slides explain them. A default is the value a reader uses when a record has no value for a field.
A data contract is an agreement between the team that produces data and the teams that consume it. A wiki page is not enough, because nobody runs it. A contract is useful when a machine checks it, and a change that breaks it is stopped before it ships.

Read the table from top to bottom. The first two rows are shape, and a schema covers them. The other three are meaning, freshness and ownership, and a schema says nothing about them. Knowing a field is a number does not tell you if it is in kilograms or pounds. Nor does it tell you if 0 is a real value or a sign that something is missing.
There is now an open standard for writing all of this down. The Open Data Contract Standard (ODCS) is a YAML format, which is plain text for settings. It is kept by Bitol, an LF AI & Data Foundation project. Its sections include the schema, data quality, a service-level agreement (SLA, the promised freshness and uptime), the team and their roles. Its current version is 3.2.0. You do not need ODCS to have a contract, but it is a good checklist of what a contract covers.
Many teams pass data through , a system that carries streams of records between programs. Each stream is called a topic. When data crosses between teams on Kafka, the shape part of the contract usually lives in a schema registry. Confluent, a company that sells a platform built on Kafka, makes the common one: Confluent Schema Registry. I describe it from its own documentation. It stores Avro, Protobuf and JSON Schema schemas. Avro and Protobuf are compact binary formats that come with a schema, and JSON Schema is a schema language for JSON.

Follow the numbered steps. The producer registers its schema, and the registry checks it against the earlier versions. If it is compatible, the registry returns an id; if not, it refuses with error 409. The producer then sends each record with that id in front, not the whole schema. The consumer reads the id, asks the registry for that schema once, keeps a copy, and decodes.
Notice what the registry never sees: the records. It checks schemas when they are registered. The part that checks each record is the producer's serializer, the code that turns a record into bytes. Confluent's Avro serializer throws an error when a record does not fit its schema. Confluent Server can also check that each message carries a registered schema id, but its documentation says this "does not perform data introspection". It does not look inside the data.
Plain JSON with no schema has no contract at all. That is how a renamed column can quietly become an empty value, as it did in lesson 2.
Schemas have to change. The question is which changes a reader on another version can survive. Confluent's documentation names three main kinds of compatibility, and the default is BACKWARD. There is also NONE, which switches the check off.

Each mode is a promise about one direction. BACKWARD promises that a reader on the new schema can read old data. Confluent's documentation then says the important part plainly: there is "no assurance that consumers using older schemas can read data produced using the new schema. Therefore, upgrade all consumers before you start producing new events." FORWARD promises the other direction, so the producers upgrade first. FULL promises both, so either side can go first.
A plain mode checks the new schema against the latest version only. A transitive mode, such as BACKWARD_TRANSITIVE, checks it against every earlier version. Confluent prefers BACKWARD for so that consumers can rewind to the start of a topic. FULL_TRANSITIVE is the strictest setting.
How does a registry decide? For Avro, it follows Avro's own rules for reading a record with a different schema from the one that wrote it.

These four rules are from the Avro specification, which calls them schema resolution. An enum is a field that may hold only one of a fixed list of words, called symbols, such as clear, misty or rain. Notice the second rule. A reader that meets a missing field uses its if it has one. This rule makes adding a field safe, and the lab shows it can also hide a missing one.
I wrote the lab's design into the docstring of its script, contract_demo.py, on 1 October 2026, before it ran. Before writing it I loaded the data once to see its columns, and timed one model fit. No score was computed before the design was written.

The data is a public table from a 2013 study by Fanaee-T and Gama: 17,379 hours of the Capital Bikeshare system in Washington, D.C., in 2011 and 2012. Each hour has twelve inputs, such as season, hour, weather, temperature, humidity and wind speed, and the number of rides. The temperatures are in what its description calls degrees Celsius. The original UCI page now gives a different scaling formula, so treat the unit as the OpenML copy's own claim. The weather field holds one of four words; the rarest, heavy_rain, appears in only 3 hours.
The producer is pretend. It sends one record per hour under a schema I wrote in Avro's format, called v1, and v1 is the contract. In v1 the weather list has three words, so the producer sends heavy rain as "rain".
The consumer reads every record with v1 and gives it to boosted trees: many small decision trees, each fixing the mistakes of the ones before. They learn from 13,003 hours, January 2011 to June 2012, and predict the other 4,376. I cut those into 30 blocks of about 146 hours each, so every result is a mean over 30 blocks with its lowest and highest. The score is the mean absolute error (MAE): how many rides an hour the prediction is off by, on average. Lower is better.
The twelve changes are the kind a producer team makes for good reasons. In every replay, the producer has deployed its change and the consumer has not upgraded yet. Right after a deploy, the old consumer is most exposed.
The first group changes the schema. Add an optional field with a default. Add a required field with no default. Delete windspeed. Rename temp to temperature. Store the hour as a long, a bigger whole number, instead of an int, a smaller one. Store humidity as a float, which keeps about 7 digits, instead of a double, which keeps about 16.
Each of these has an ordinary reason behind it. A new field holds a new signal the producer team wants to share. A field is deleted because the sensor behind it was retired. A rename makes a name clearer. A bigger number type is room to grow, and a smaller one saves storage. None of them is careless. In the lab, every one of the changes that hurt the consumer is the kind of change a careful team makes without knowing who reads the data.
The second group keeps the schema and changes the data. These happen when the producer switches to a new source, such as a weather service that reports in Fahrenheit. Send temperatures in Fahrenheit, with the same names and types. Send humidity as a percent, 0 to 100, instead of 0 to 1. Then two enum changes: add heavy_rain to the weather list, and change the order of the season list.
The last two repeat a delete and a rename, but in schemas where every number field has a default of 0. Teams add defaults like these on purpose, because under FULL a field can only be deleted later if it has one.
For each change, my own checker gives the verdict of each mode. It is a few dozen lines following the Avro rules. Then the old consumer reads all 17,379 records the producer would send, and I count the ones it cannot decode. When every serving record decodes, I score the trees on what the consumer actually read.
Here are the verdicts. I also asked the official Apache Avro library, version 1.12.2, the same questions in a separate small program. Its verdicts and its failed-record counts matched mine on all twelve changes.

The table has one row per change. Read the "failed" column first: it is the old consumer's view. Three changes made every single record unreadable, 17,379 of 17,379. Adding heavy_rain failed on only 3 records, the three heavy rain hours.
The verdicts follow the rules from the last slide, so they are no surprise once you know them.
Take the required field. Under BACKWARD, a new reader reads old data, which has no station_id, and the field has no default, so it fails. Under FORWARD, an old reader reads new data and simply skips station_id, so it passes. Humidity as a float is the mirror image. A new reader expecting a float cannot read an old double, so BACKWARD fails. An old reader expecting a double can read a float, so FORWARD passes. Each mode looks in one direction only, and FULL looks in both. The interesting part is how they line up with the "failed" column, which the next figure counts.

BACKWARD, the default, passed 9 of the 12 changes. FORWARD passed 8 and FULL 6. The block heights are those counts. Passing more is not bad in itself. A mode that passes a change has done its job if the reader it protects can read the data. What matters is which reader that is, and whether the data still means the same thing. The caption holds the two results that matter, and the next slides take them one at a time.
Four changes left the old consumer unable to read some or all records. BACKWARD passed three of them. FORWARD and FULL passed none.

Each box down the chart is one step of what happened. Deleting windspeed passed BACKWARD, because a new reader can simply skip it in old data. But the old reader still needs windspeed, and v1 gave it no default. So all 17,379 records failed. Storing the hour as a long also passed BACKWARD, because a new reader can read an old int as a long. The old reader cannot read a long as an int, so every record failed again.
The heavy_rain change failed the same way, on 3 records. The old weather list had no heavy_rain and no default. Avro lets an enum carry a default symbol for exactly this case. But think about what it would do here: an hour of heavy rain would arrive as, say, "clear". The failure would turn into a quiet wrong value, the same trade the defaults make on the next slides.
This is not a bug in BACKWARD. It is exactly what its documentation says it does not promise. Under BACKWARD the consumers must upgrade first. If your producer team deploys first, BACKWARD makes no promise to your consumers. Here it stopped 1 of the 4 changes that broke them: the rename with no defaults. It did so only because that change also fails in the direction BACKWARD does check.
One more thing. Upgrading the consumer first fixes the decoding, but not the model. A consumer on the new schema can read records with no windspeed, yet its trees were trained with windspeed. Compatibility is about reading records, not about whether your model still has its inputs.
Now the quiet ones. Nine changes left every serving record readable. For each, I scored the trees on what the old consumer read and compared that with the clean run. The clean MAE was 44.5 rides an hour, from 30.7 to 66.3 across the 30 blocks.

Each dot is one block. Sending temperatures in Fahrenheit raised the error by 30.20 rides an hour on average, and in 28 of 30 blocks; it ranged from -1.9 to 79.5. Humidity as a percent raised it by 30.41, in 27 of 30 blocks. The rename with defaults of 0 raised it by 36.39, in 27 of 30 blocks.
Deleting windspeed with a default of 0 raised it by only 0.27, and in just 16 of 30 blocks, from -0.8 to 2.3. That is not a clear loss of accuracy. One possible reason, which I checked after the results: windspeed was already exactly 0 in 1,529 training hours, 11.8 percent. The trees had seen plenty of zeros.

The bars count something simpler: the share of the 4,376 serving hours whose prediction moved by more than one ride. Fahrenheit moved 99.3 percent of them. Even the windspeed default moved 36.1 percent, though the error barely changed. Its answers changed, without clearly getting worse.
Every one of these four passed every mode, FULL included. This is the lab's headline.
The rename shows the trap best, because I ran it twice.

With no defaults, every mode rejected the rename. Had someone shipped it anyway, the old consumer could not have read a single record. That is a loud failure, and loud failures get fixed.
With defaults of 0 on both sides, every mode passed it. To the registry, a rename is one field deleted and another added. The new reader fills the new name from its default when reading old data. The old reader fills temp from its default when reading new data. So both directions work, and the old consumer read a temperature of 0 for every hour. Nothing failed. The error rose by 36.39 rides an hour.
The earlier version of this lesson had an Avro example with "default": 0 on a person's age. This is the same trap: an age of 0 is a real-looking value, so a missing age would look like a newborn.

The float change is a smaller lesson, found after the results. A float keeps fewer digits, so 0.81 arrives as 0.8100000023841858. That changed exactly one prediction in 4,376, by 7.080 rides. One possible reason, a guess: that hour's humidity sat right on one of the trees' cut points. Tiny changes are not always invisible.
A shape check cannot see units. A check on the data's meaning can. I call it a semantic contract: rules about what the values mean, checked on the data at run time. Lesson 5 measured checks like these on their own, including their false alarms. Here the question is narrower: did they catch the four changes the schema check let through?

The rules came from the training hours only, and I wrote them into the design before the run. I must be honest about one thing. I wrote them knowing the twelve changes, so a rule that catches a change is partly right by design. The test is weaker than it looks, and I say so. What the run adds is that the rules fired on 0 of the 30 clean blocks.

Each of the four changes set off a rule in all 30 blocks. The range rules caught the units. The rule that a field must take at least two values in a block caught both defaults, because a default of 0 every hour is a constant. A range rule alone would have missed the rename: 0 is inside the temperature range.

After the results I looked at how close the Fahrenheit case came to slipping through. In the late autumn blocks, the coldest hours read about 42 to 50, inside the range. Only the warmer hours of each block crossed 50.18. The coolest block's hottest hour read 57.1. In a colder city or a colder month, a range rule could miss a unit change for days.
A plain mode compares a new schema with the latest one only. That leaves a gap, and the lab has one chain of changes to show it.

Version 2 deletes windspeed, which BACKWARD allows. Version 3 adds it back as text, with a default, which BACKWARD also allows against version 2. But a topic may still hold old records from version 1, where windspeed was a number. A version 3 reader cannot read a number as text. BACKWARD never compared version 3 with version 1, so it never saw this. BACKWARD_TRANSITIVE does. Confluent's documentation links a similar example, where removing a default is the step that breaks the chain.

Adding heavy_rain shows a different kind of delay. Only 3 of 17,379 hours had heavy rain, all in January. None fell in the half I scored. So a consumer that could not read the new value would have run for months without a single error, then failed on the first hour of heavy rain. A rare new value is a break that waits.
A contract that nobody checks is a suggestion. The lab points to two checks, and each catches what the other misses.
At the merge, in CI. CI, continuous integration, is the set of automatic checks that run on every proposed code change. Keep the schema files in git next to the producer's code. When a pull request edits one, the build asks the registry's compatibility API whether the new schema passes the topic's mode, and fails if it does not. The rename with no defaults would die here, before anything shipped.
On the data, at run time. The checks on meaning from two slides back run on the data as it arrives. A failing batch is held back, and a person is told. Lesson 5 covers how to do this with Great Expectations and Pandera, and how to set the rules so they do not fire all day. The two tools differ in one way. Great Expectations returns a result with success set to false, and your pipeline must act on it. Pandera raises an error by default.
If your contract lives in a warehouse rather than on , dbt has its own version of the first check. A warehouse is a database built for analysis, and dbt is a tool that builds its tables from files. A dbt model contract (dbt 1.5 and later) checks a model's column names and types before it is built. Which constraints it actually enforces, such as not null, depends on the database.
Some changes cannot be made compatible, and some should not be. A field that has meant the wrong thing for a year needs a new name, not a quiet fix. Those changes need a plan instead of a check.

The steps go from top to bottom. Confluent's documentation describes the choice plainly: upgrade every producer and consumer at the same time, or "more likely" create a brand-new topic and move applications to it. Writing both old and new for a while lets each consumer move on its own schedule. The owner list in the contract tells you who to tell.
Confluent also sells a newer way, for its Enterprise and Cloud plans. Its data contracts can group schemas by a major version and attach migration rules that translate records between versions. Check your own plan and version before you count on it.
Who owns the contract? The producer, because only the producer can keep the promise. But the consumers should review changes. On GitHub, a CODEOWNERS file asks the right people to review a change to the schema file. It only blocks the merge if the branch rule "Require review from Code Owners" is turned on.
Google's Rules of Machine Learning define training-serving skew as "a difference between performance during training and performance during serving". It lists three causes. The first is "a discrepancy between how you handle data in the training and serving pipelines". The second is a change in the data between training and serving, and the third a feedback loop.
The lab is the second cause, made small. The trees learned on temperatures in what the data calls Celsius, and were served temperatures in Fahrenheit. Nothing in the model was wrong. The data changed its meaning between training and serving, and the error rose by 30.20 rides an hour.
So a contract does not remove skew. It removes some of its causes, and only where it is checked. The registry and the serializer protect the shape. The checks on meaning catch some changes of meaning. A is a service that keeps a model's inputs, called features, for both training and serving. Feast, an open-source one, gives training and serving the same feature definitions, though Feast does not run the pipelines that compute them. None of these touches a feedback loop.
There is a practical point for ML teams here. A data engineer watches for jobs that fail. A model owner has to watch for jobs that succeed with the wrong data, because those are the ones that move predictions. In the lab, the four changes that hurt the model were exactly the four that raised no error. So the checks on meaning belong to whoever owns the model, not only to whoever owns the pipeline.
Data contracts became a popular idea around 2022. Here is what the sources I checked actually say, and no more.

Convoy, a freight company, published one of the most cited guides in 2022. Its engineer Adrian Kreuziger wrote that "data contracts must be enforced at the producer level", using Protobuf or Avro, checked in CI. Because Convoy's producers update first, he set the registry to FORWARD, not the BACKWARD default. Convoy shut down its operations in October 2023.
That choice matches the lab. FORWARD looks from the old reader's side, so it let through none of the 4 changes that broke the old consumer here. It still let through all 4 quiet ones, because those were about meaning, which no mode checks. If your producers deploy first, FORWARD protects the readers who have not caught up yet.
Airbnb wrote in 2020 about rebuilding its data quality. The earlier version of this lesson said that work was about unannounced upstream changes. The posts do not say that. They say data ownership responsibilities "were not clearly defined", and that this "was a bottleneck when issues arose".

The cards list the tools from their own documentation. One more fact from them is worth knowing. With Protobuf, field numbers travel in the data, not names, so renaming a field does not break old binary readers. It does break any reader that uses JSON or the old name in its code.
This script is the lab. It downloads the data, replays the twelve changes and the chain, and prints what each mode did and what the model saw. It needs no GPU. On my laptop it ran in a few seconds, and it prints no timings.

Before you run this lab. You need Python 3 and two libraries: pip install scikit-learn pandas. scikit-learn does the download and the model, and brings NumPy with it. pandas is what scikit-learn uses to read the downloaded table. The first run downloads the data from OpenML (about 0.2 MB), so it needs an internet connection once. After that, scikit-learn keeps a copy in a folder in your home directory (scikit_learn_data).
I ran it with scikit-learn 1.9.1 on a Mac. These libraries run on Windows and Linux too, but I have not checked the numbers there. Give it a file name, python contract_demo.py out.json, and it also saves every number. That is how results/ctr-demo.json was made. Two runs gave byte-for-byte the same file. One change came after the first run, and the docstring says so. The saved file now also holds each change's two schemas, so the Avro check can read them. Nothing printed changed.
"""Data contracts: which schema changes does a registry let through,
and what does each one do to a model that reads the data?
Lesson 10 of 'Data Engineering for ML', made small. It needs Python 3
with scikit-learn and pandas (pip install scikit-learn pandas); NumPy
comes with scikit-learn. The first run downloads the bike share data
from OpenML (about 0.2 MB) and keeps a copy. After that it runs in
well under a minute on a laptop.
python contract_demo.py # print the results
python contract_demo.py out.json # and save every number
Design, written 2026-10-01 before the first run. Before writing it I
loaded the data once to see its columns, and timed one model fit. No
score was computed before this was written.
Data: OpenML 42712, 'Bike_Sharing_Demand': 17,379 hours of a bike
share, 2011 and 2012. Twelve input columns and the rides that hour.
The producer sends one record per hour. Its v1 schema, in Avro's
JSON form, is the contract the consumer was built on. v1 has three
weather values, so the producer sends heavy rain as 'rain'.
The consumer reads every record with the v1 schema and gives it to
boosted trees (HistGradientBoostingRegressor, early stopping off,
random_state 0), trained on Jan 2011 to Jun 2012. It serves Jul to
Dec 2012, cut into 30 blocks of consecutive hours.
The checker is my own code for the Avro specification's 'schema
resolution' rules: a reader field is matched by name or alias; a
field only the writer has is skipped; a field only the reader has
takes the reader's default, else it is an error; int may be read as
long, float or double, long as float or double, float as double; an
enum symbol the reader lacks is an error unless the reader's enum
has a default; a union needs a branch that can read the other side.
Modes, as in Confluent Schema Registry: BACKWARD, the new schema can
read data written with the old one; FORWARD, the old can read the
new; FULL, both. Plain modes check against the latest version only;
BACKWARD_TRANSITIVE checks against every earlier version.
Twelve producer changes, each a v2 schema plus what the producer
now sends: 1 add an optional field; 2 add a required field with no
default; 3 delete windspeed; 4 rename temp to temperature; 5 widen
hour from int to long; 6 narrow humidity from double to float; 7
temp and feel_temp in Fahrenheit, same names, same types; 8 humidity
as a percent, same name, same type; 9 add heavy_rain to the weather
enum; 10 reorder the season enum; 11 delete windspeed when every
number field has a default of 0; 12 rename temp the same way.
And one chain: v2 deletes windspeed, v3 adds it back as a string.
For each change: the verdict of each mode against v1. Then the
consumer, still on v1, reads all 17,379 records sent after the
change: how many fail to decode. When every serving record decodes:
the share of the 4,376 serving hours whose prediction moved by more
than one ride, and the mean absolute error (MAE, rides per hour) of
each block against the clean run: mean, lowest and highest over the
30 blocks, and the blocks where the change raised it.
A semantic contract, set from the training rows only: temp and
feel_temp within 10 of their training range; humidity from 0 to 1;
windspeed from 0 to 10 above its training maximum; every number
field takes at least two values in a block. For each change, and
for the clean run, the blocks where a rule fires.
No timings are printed.
Changed after the first run: the saved file now also holds each
change's two schemas, so another program can check them. Nothing
printed changed. Then one fix, after a review: the checker crashed
on a change from an enum or record to a plain type (such as season
sent as a string), and did not compare names. It now refuses both
cases, as Avro does. The saved file and the printout did not change.
Author: Roni Das
Created: 2026-10-01
"""
import json
import sys
import numpy as np
import sklearn
from sklearn.datasets import fetch_openml
from sklearn.ensemble import HistGradientBoostingRegressor
# ── the producer's v1 schema: the contract ──
NUMS = ["temp", "feel_temp", "humidity", "windspeed"]
SEASONS = ["spring", "summer", "fall", "winter"]
WEATHER = ["clear", "misty", "rain"]
def field(name, type_, default=None):
f = {"name": name, "type": type_}
if default is not None:
f["default"] = default
return f
def enum(name, symbols):
return {"type": "enum", "name": name, "symbols": symbols}
def record(fields):
return {"type": "record", "name": "Hour", "fields": fields}
V1_FIELDS = [field("season", enum("Season", SEASONS)),
field("year", "int"), field("month", "int"),
field("hour", "int"), field("holiday", "boolean"),
field("weekday", "int"), field("workingday", "boolean"),
field("weather", enum("Weather", WEATHER))]
V1_FIELDS += [field(n, "double") for n in NUMS]
V1 = record(V1_FIELDS)
def with_defaults(fields):
"""The same fields, every number field given a default of 0."""
return [dict(f, default=0.0) if f["type"] == "double" else f
for f in fields]
V1D = record(with_defaults(V1_FIELDS))
# ── the checker: Avro's schema resolution rules ──
PROMOTE = {"int": {"long", "float", "double"},
"long": {"float", "double"}, "float": {"double"},
"string": {"bytes"}, "bytes": {"string"}}
def can_read(reader, writer):
"""Can data written with `writer` be read with `reader`?"""
if isinstance(writer, list): # a writer union
return all(can_read(reader, w) for w in writer)
if isinstance(reader, list): # a reader union
return any(can_read(r, writer) for r in reader)
if isinstance(reader, str) and isinstance(writer, str):
return reader == writer or reader in PROMOTE.get(writer, ())
if isinstance(reader, str) or isinstance(writer, str):
return False # a plain type against an enum or record
if (reader["type"], reader["name"]) != (writer["type"], writer["name"]):
return False # Avro: named types must match
if reader["type"] == "enum":
extra = set(writer["symbols"]) - set(reader["symbols"])
return not extra or "default" in reader
wf = {f["name"]: f for f in writer["fields"]}
for f in reader["fields"]:
names = [f["name"]] + f.get("aliases", [])
match = [wf[n] for n in names if n in wf]
if match:
if not can_read(f["type"], match[0]["type"]):
return False
elif "default" not in f:
return False
return True
def verdicts(old, new):
back, fwd = can_read(new, old), can_read(old, new)
return {"BACKWARD": back, "FORWARD": fwd, "FULL": back and fwd}
def decode(datum, writer, reader):
"""One record written with `writer`, read with `reader`."""
if isinstance(reader, str) or isinstance(writer, str):
if not can_read(reader, writer):
raise ValueError(f"{writer} as {reader}")
return float(datum) if reader in ("float", "double") else datum
if (reader["type"], reader["name"]) != (writer["type"], writer["name"]):
raise ValueError("a different type or name")
if reader["type"] == "enum":
if datum not in reader["symbols"]:
raise ValueError(f"unknown symbol {datum}")
return datum
wf = {f["name"]: f for f in writer["fields"]}
out = {}
for f in reader["fields"]:
names = [f["name"]] + f.get("aliases", [])
hit = [n for n in names if n in wf]
if hit:
out[f["name"]] = decode(datum[hit[0]], wf[hit[0]]["type"],
f["type"])
elif "default" in f:
out[f["name"]] = f["default"]
else:
raise ValueError(f"no value and no default: {f['name']}")
return out
# ── the twelve changes: a new schema, and what the producer sends ──
def swap(fields, name, new):
return [new if f["name"] == name else f for f in fields]
def drop(fields, name):
return [f for f in fields if f["name"] != name]
def renamed(r, old, new):
return {(new if k == old else k): v for k, v in r.items()}
def f32(x):
return float(np.float32(x))
def to_f(c):
return c * 9 / 5 + 32
CHANGES = {
"add optional field": (
V1, record(V1_FIELDS + [{"name": "uv_index",
"type": ["null", "double"],
"default": None}]),
lambda r: r | {"uv_index": None}),
"add required field": (
V1, record(V1_FIELDS + [field("station_id", "string")]),
lambda r: r | {"station_id": "dc-1"}),
"delete windspeed": (
V1, record(drop(V1_FIELDS, "windspeed")),
lambda r: {k: v for k, v in r.items() if k != "windspeed"}),
"rename temp": (
V1, record(swap(V1_FIELDS, "temp",
field("temperature", "double"))),
lambda r: renamed(r, "temp", "temperature")),
"hour int to long": (
V1, record(swap(V1_FIELDS, "hour", field("hour", "long"))), None),
"humidity to float": (
V1, record(swap(V1_FIELDS, "humidity",
field("humidity", "float"))),
lambda r: r | {"humidity": f32(r["humidity"])}),
"temps in Fahrenheit": (
V1, V1, lambda r: r | {"temp": to_f(r["temp"]),
"feel_temp": to_f(r["feel_temp"])}),
"humidity as percent": (
V1, V1, lambda r: r | {"humidity": r["humidity"] * 100}),
"add enum heavy_rain": (
V1, record(swap(V1_FIELDS, "weather", field(
"weather", enum("Weather", WEATHER + ["heavy_rain"])))),
"heavy"),
"reorder season enum": (
V1, record(swap(V1_FIELDS, "season", field(
"season", enum("Season", SEASONS[::-1])))), None),
"delete windspeed, defaults": (
V1D, record(drop(with_defaults(V1_FIELDS), "windspeed")),
lambda r: {k: v for k, v in r.items() if k != "windspeed"}),
"rename temp, defaults": (
V1D, record(swap(with_defaults(V1_FIELDS), "temp",
field("temperature", "double", 0.0))),
lambda r: renamed(r, "temp", "temperature")),
}
# ── the data ──
d = fetch_openml(data_id=42712, as_frame=True, parser="auto")
df = d.frame
truth_weather = df["weather"].astype(str).tolist()
rows = []
for t in df.itertuples(index=False):
rows.append({"season": str(t.season), "year": int(t.year),
"month": int(t.month), "hour": int(t.hour),
"holiday": str(t.holiday) == "True",
"weekday": int(t.weekday),
"workingday": str(t.workingday) == "True",
"weather": "rain" if t.weather == "heavy_rain"
else str(t.weather),
"temp": float(t.temp), "feel_temp": float(t.feel_temp),
"humidity": float(t.humidity),
"windspeed": float(t.windspeed)})
y = df["count"].to_numpy(float)
train = ((df["year"] == 0) | (df["month"] <= 6)).to_numpy()
serve = np.flatnonzero(~train)
blocks = np.array_split(serve, 30)
def features(recs):
"""The consumer's own encoding of v1 records, by symbol name."""
return np.array([[SEASONS.index(r["season"]), r["year"], r["month"],
r["hour"], r["holiday"], r["weekday"],
r["workingday"], WEATHER.index(r["weather"])]
+ [r[n] for n in NUMS] for r in recs], float)
X = features(rows)
model = HistGradientBoostingRegressor(early_stopping=False,
random_state=0)
model.fit(X[train], y[train])
clean = model.predict(X[serve])
pos = {i: k for k, i in enumerate(serve)}
def block_mae(pred):
return [float(np.abs(pred[[pos[i] for i in b]] - y[b]).mean())
for b in blocks]
# the semantic contract, set from the training rows only
Xt = X[train]
LIMITS = {"temp": (Xt[:, 8].min() - 10, Xt[:, 8].max() + 10),
"feel_temp": (Xt[:, 9].min() - 10, Xt[:, 9].max() + 10),
"humidity": (0.0, 1.0),
"windspeed": (0.0, Xt[:, 11].max() + 10)}
def semantic(Xs):
"""Per block, the rules that fire on the decoded serving rows."""
fired = []
for b in blocks:
rule = []
B = Xs[[pos[i] for i in b]]
for j, n in enumerate(NUMS):
col = B[:, 8 + j]
lo, hi = LIMITS[n]
if col.min() < lo or col.max() > hi:
rule.append(f"{n} range")
if len(np.unique(col)) < 2:
rule.append(f"{n} constant")
fired.append(rule)
return fired
print(f"scikit-learn {sklearn.__version__}")
print(f"hours {len(rows):,}; train {train.sum():,}; "
f"serve {len(serve):,} in {len(blocks)} blocks")
print(f"heavy rain hours: {truth_weather.count('heavy_rain')}")
base = block_mae(clean)
out = {"version": sklearn.__version__, "hours": len(rows),
"train": int(train.sum()), "serve": int(len(serve)),
"heavy_rain": truth_weather.count("heavy_rain"),
"limits": {k: [float(v) for v in lim] for k, lim in LIMITS.items()},
"clean": {"block_mae": base,
"semantic": semantic(X[serve])},
"changes": {}}
print(f"clean MAE {np.mean(base):.1f} rides/hour "
f"({min(base):.1f} to {max(base):.1f})")
print(f"clean blocks a semantic rule fired on: "
f"{sum(bool(f) for f in out['clean']['semantic'])}")
print("\nverdicts against v1 (B=BACKWARD F=FORWARD L=FULL)")
print(f" {'change':<28}{'B':>2}{'F':>2}{'L':>2}{'errors':>8}")
for name, (old, new, send) in CHANGES.items():
v = verdicts(old, new)
sent = []
for r, w in zip(rows, truth_weather):
if send == "heavy":
sent.append(r | {"weather": w})
else:
sent.append(send(r) if send else r)
got, errors = [], 0
for s in sent:
try:
got.append(decode(s, new, old))
except ValueError:
got.append(None)
errors += 1
res = {"verdicts": v, "errors": errors,
"schemas": {"old": old, "new": new}}
if all(got[i] is not None for i in serve):
Xs = features([got[i] for i in serve])
pred = model.predict(Xs)
mae = block_mae(pred)
diff = [m - c for m, c in zip(mae, base)]
res |= {"moved": float(np.mean(np.abs(pred - clean) > 1)),
"block_mae": mae,
"mae_change": {"mean": float(np.mean(diff)),
"min": float(min(diff)),
"max": float(max(diff))},
"worse_blocks": int(sum(x > 0 for x in diff)),
"semantic": semantic(Xs)}
out["changes"][name] = res
mark = "".join(f"{'y' if v[m] else '-':>2}"
for m in ("BACKWARD", "FORWARD", "FULL"))
print(f" {name:<28}{mark}{errors:>8,}")
print("\nwhat the consumer's model saw (serving hours)")
print(f" {'change':<28}{'moved':>7}{'MAE +':>8}{'rules':>7}")
for name, res in out["changes"].items():
if "moved" not in res:
continue
fired = sum(bool(f) for f in res["semantic"])
print(f" {name:<28}{res['moved']:7.1%}"
f"{res['mae_change']['mean']:8.2f}{fired:>4}/30")
# the chain: each step passes BACKWARD, the whole does not
V2 = record(drop(V1_FIELDS, "windspeed"))
V3 = record(drop(V1_FIELDS, "windspeed")
+ [field("windspeed", "string", "")])
chain = {"v2 after v1": can_read(V2, V1), "v3 after v2": can_read(V3, V2),
"v3 after v1": can_read(V3, V1)}
out["chain"] = chain | {"schemas": [V1, V2, V3]}
print("\nchain: v2 drops windspeed, v3 adds it as a string")
for k, ok in chain.items():
print(f" BACKWARD {k}: {'passes' if ok else 'fails'}")
if len(sys.argv) > 1:
json.dump(out, open(sys.argv[1], "w"), indent=1)

The report lives in scripts/labs/dataeng/contract_report.py. It does not trust the demo's code for anything it can redo another way. It reads the downloaded data file itself, line by line. It decodes every record with its own resolver, written apart from the demo's checker. Then it hashes the result exactly as the Avro check does. A hash is a short fingerprint of the data that changes if any value changes.
The Avro check, contract_avro_check.py, runs in its own small environment with the official avro package, version 1.12.2. It writes every record in Avro's real binary format under the new schema and reads it back under the old one. The report's hashes must equal Avro's for all twelve changes, so what the model was fed is what real Avro decoding gives.
The model must be the same model, so the report refits it with scikit-learn. Everything after the fit is its own code: the blocks, the errors, the moved shares and the rules. It stops unless all 444 checks agree.
What came when: the demo's design came before any run. The rule-by-rule counts came after the results, and so did four other facts: the 32-bit float, the windspeed zeros, the heavy rain hours and the Fahrenheit blocks. The report's docstring lists them, and so does the Avro probe below.

The checker and the decoder are plain Python. scikit-learn downloads the data and fits the trees, NumPy holds the arithmetic, and pandas reads the downloaded table. Apache Avro gave the second opinion. After the results I also asked it about two features the Avro specification lets a library leave out. One is a reader alias, which renames a field. The other is an enum default. Its checker accepted both. Its Python reader applied neither and raised an error on both. An alias is only as good as the library that reads it, so test yours.
This box holds the same checker as the lab, the twelve changes, and what the lab measured for each. It has no model and no data rows. It runs in your browser.
As it is, the box checks the Fahrenheit change. All three modes pass it, because the schema did not change at all. Then it prints what the lab measured: no failed records, 99.34 percent of predictions moved, and the error up 30.20 rides an hour.
Now set CHANGE = 'delete windspeed'. BACKWARD passes, FORWARD and FULL refuse, and the measured line says all 17,379 records failed. Then try 'rename temp, defaults' and watch every mode pass a change that fed the model a temperature of 0. To go further, write your own change into CHANGES and see what the checker says.
The v1 schema. field, enum and record build Avro schemas as Python dictionaries, in Avro's own JSON shape. V1_FIELDS is the contract: twelve fields, two enums and four numbers. with_defaults gives every number field a default of 0.
can_read. This is the checker, in about twenty lines. It walks the two schemas and applies Avro's rules. A union is a field that may hold one of several types, such as nothing or a number, and it needs a branch that fits. A plain type must match or be promotable, which means readable as a wider type. An enum must know every symbol unless it has a default. A record needs a value or a default for each of the reader's fields. verdicts calls it in both directions to get BACKWARD, FORWARD and FULL.
decode. The same rules, applied to one record. It returns what the old consumer would read, or raises an error. A default fills a missing field here, which is how the rename with defaults turned into zeros.
CHANGES. Each change is three things: the old schema, the new schema, and a small function for what the producer now sends.
The main run. fetch_openml gets the data, and each row becomes a v1 record. The trees train once, on the clean training hours. For each change, the script checks the verdicts, decodes all 17,379 records with , and, if every serving record decoded, predicts and scores each block. applies the four rules. Last comes the chain, three calls to .

Write the contract down. Include the units, the allowed values, how fresh the data will be, and who owns it. A schema alone covers only the shape.
Pick a mode that matches who upgrades first. If your consumers upgrade first, BACKWARD fits; if producers do, FORWARD. If you cannot control the order, use FULL. Consider a transitive mode if old data stays on the topic.
Check every schema change in CI. Ask the registry's compatibility API before the merge, not after the deploy.
Be careful with defaults. Choose defaults that cannot pass for real values, or none, and know which fields have them. A default of 0 hid a missing temperature in every hour of the lab.
Check the meaning at run time. Ranges in the right units, allowed values, and a rule against a field that never changes. Here those caught all four quiet changes.
Plan the changes no mode allows. A new version, both written for a while, consumers moved one by one, then the old one retired.

A full contract is worth it when another team owns the data your model reads, or many consumers read the same topic. It is also worth it when a wrong value costs more than a stopped job. The lab's quiet changes are the reason. In the lab, the model served half a year of hours with temperatures in the wrong unit, and no error was raised anywhere.
A lighter version is fine when one team owns both ends and ships them together, or the data is a one-off file for one analysis. It is also fine when a person reads every output before anyone uses it. Then a schema check and a few tests may be enough, because the person is the check on meaning.

One dataset and one model. Another model might lean less on temperature, and so suffer less from a unit change. The sizes of the errors are this model's, on this data.
Avro only. Protobuf and JSON Schema have their own rules. A Protobuf rename, for example, does not break old binary readers.
The verdicts are rules, not findings. Given the schemas, Avro's rules decide the verdicts, and Apache Avro agreed with my checker. What the lab measured is what each change did to the consumer and its model.
Rules written knowing the changes. The meaning rules caught all four quiet changes, but I chose both the rules and the changes. Lesson 5 measured how often checks like these fire on good data.
What came after the results. The rule-by-rule counts, the 32-bit float, the windspeed zeros, the heavy rain hours, the Fahrenheit blocks and the Avro alias probe.

If your model reads data another team writes, ask the first two questions this week. Is the contract written down with units and an owner? Which mode is the registry set to, and does it match who deploys first?
Then look at the defaults. For each field with a default, ask what a missing value would become, and whether your model could tell it from a real one.

The card keeps the lesson's three numbers. A rename that every mode passed raised the error by 36.39 rides an hour. BACKWARD passed 9 of 12 changes, and 3 of those left the old consumer unable to read its records. Rules about meaning fired on all four quiet changes, in every block, though I wrote those rules knowing the changes.
The next lesson in this chapter is about label noise: what wrong labels in the training data do to a model, measured.
4 questions - Score 80% to pass
Under BACKWARD, the registry's default, which reader does the mode make no promise to?
The producer started sending temperatures in Fahrenheit, with the same field names and types. What did the lab measure?
Both schemas gave every number field a default of 0, and then temp was renamed. Why was this worse than the rename with no defaults?
Each step of the windspeed chain passed BACKWARD: v2 deleted it, v3 added it back as text. What does BACKWARD_TRANSITIVE add?
This is a real run in VS Code's terminal (python contract_demo.py).

When I ran it, it printed the version, the counts, the verdicts, what the model saw, and the chain. All of it matches the stored ctr-demo.json, and the longest printed line was 52 characters. The report's demo mode checks those lines against the file.
To try something I have not run, add a thirteenth change to CHANGES: send windspeed in miles per hour instead of its current unit. I cannot tell you what it prints, because I have not run it. The question to ask is whether any of the four meaning rules notices.
decodesemanticcan_read