Data Engineering For Ml

Data Contracts and Schema Management for ML Pipelines

0 of 27 complete

0%

Contents

Back|Data Engineering For MlData Contracts and Schema Management for ML Pipelines
1/27
66 min left
  1. Home
  2. AI Engineering: Data, RAG and Agents
  3. Data Engineering for ML
  4. Data Contracts and Schema Management for ML Pipelines
Prerequisites
Synthetic Data Generation: When Fake Data Beats No Datarequired
Related Topics
Fine-Tuning vs RAG vs Prompting: Choosing Your ApproachLLM and GenAI OpsParameter-Efficient Fine-Tuning: LoRA and QLoRALLM and GenAI OpsEvaluating LLMs in Production: Grading Answers That Have No Right AnswerLLM and GenAI OpsPrompt Management and Versioning: Treat Prompts as Production CodeLLM and GenAI OpsVector Databases and Approximate Nearest Neighbor SearchLLM and GenAI Ops
1 of 27
Previous lesson
Synthetic Data Generation: When Fake Data Beats No Data
Next lessonWrong Labels: How Many Can a Model Survive, and Can You Find Them?

System Design

  • Foundation
  • Intermediate
  • Advanced
  • Capstone

AI Engineering

  • Foundation
  • Data, RAG and Agents
  • Evaluation, LLM Ops and Security

systemdesign.academy

  • Home
  • Glossary
  • Interview prep
  • Reviews
  • About
  • Privacy
  • Terms

The Delivery Note

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.

An illustration of three colleagues around a table, beside text. Headed a supplier, a kitchen and a delivery note, titled the note still says 25; it means something else. Beside the picture: a kitchen and its supplier agree what each line of the delivery note means; the kitchen's checks look at the note's shape, a name, then a number. Beneath: one day the supplier starts counting flour in pounds, not kilograms, and the note still says 25; the shape is the same, so every check passes, and the bread comes out wrong. Last line: this lesson measures which changes a shape check stops, and which ones reach the model.

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.

Where This Lesson Starts

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.

A flowchart headed the plan of this lesson, titled from the promise to the measurement. What a data contract is, and where it lives, leads to the rules: what each compatibility mode allows. That leads to the lab: twelve real producer changes, replayed. Then what each mode let through, and what reached the model. Then what a shape check cannot see, and how to change a schema safely. Beneath: measured on 17,379 real hours of a bike share.

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.

Words for This Lesson

A hand-drawn list headed eleven words for this lesson, titled the words the numbers need. Producer: the team, and the code, that writes the data. Consumer: a team or job that reads it, such as a model's feature job. Schema: the written shape of a record, field names and types. Data contract: the producer's written promise, shape, meaning, freshness and who owns it. Schema registry: a service that stores every version of a schema and checks each new one. Compatibility: whether a reader on one version can read data written with another. BACKWARD: the new schema can read data written with the old one. FORWARD: the old schema can read data written with the new one. FULL: both at once. Transitive: checked against every earlier version, not just the last. Default: the value a reader uses when a record has no value for a field. Last line: the default matters more than it looks.

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.

What a Data Contract Pins Down

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.

A two-column table headed what a data contract pins down, titled shape on one side, meaning on the other. Left, the part of the contract; right, what can check it. Field names and their structure: the registry, before a schema is accepted. The type of each field: the serializer, when a record is written. Meaning, units, allowed values and ranges: a check on the data itself, at run time. Freshness and quality promises: a check on the data, over time. An owner who answers when it breaks: a person, named in the contract. Beneath, left: a schema covers the first two. Right: the rest needs checks and people.

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.

Schema Registries: Where the Shape Is Checked

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.

A sequence diagram with four columns, producer, registry, Kafka topic and consumer, headed Confluent Schema Registry, one schema and one record, titled the registry checks schemas, not messages. Step 1, producer to registry: register this schema. Step 2, registry to producer: compatible? an id; if not, 409. Step 3, producer to Kafka topic: record, with the id in front. Step 4, topic to consumer: the record. Step 5, consumer to registry: the schema for this id. Step 6, the consumer to itself: decode, then cache it. Beneath: a record that does not fit its schema fails in the producer's serializer, before step 3.

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.

The Rules a Registry Checks

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.

A page in four labelled zones, headed compatibility types, in Confluent's own words, titled each mode makes a different promise. BACKWARD, the default: consumers using the new schema can read data written with the old one; upgrade all consumers before producing new events. FORWARD: consumers using the old schema can read data written with the new one; upgrade the producers first. FULL: both directions; producers and consumers can be upgraded independently. _TRANSITIVE: the same promise, checked against every earlier version, not only the latest one. Beneath: the default is plain BACKWARD; FULL_TRANSITIVE is the strictest.

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.

A hand-sketched column of four boxes, headed Avro's schema resolution rules, which the lab's checker follows, titled how a reader reads a record it did not write. A field only the writer has: the reader skips it. A field only the reader has: the reader's default, or an error if there is none. int can be read as long, float or double; long cannot be read as int. An enum symbol the reader lacks: an error, unless its enum has a default. Beneath: fields are matched by name, so their order does not matter.

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.

The Lab: One Producer, One Consumer, Twelve Changes

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.

A page in four labelled zones, headed what the lab ran: contract_demo.py, designed before it ran, titled one producer, one consumer, twelve changes. The data: OpenML 42712, 17,379 hours of a Washington, D.C. bike share, 2011 and 2012; twelve inputs and the rides. The producer: sends one record per hour under a v1 Avro schema, the contract. The consumer: reads with v1, feeds boosted trees trained on 13,003 hours; serves 4,376 hours in 30 blocks. The changes: twelve v2 schemas, each with what the producer now sends; the consumer has not upgraded. Beneath: no timings are printed; the rule counts came after the results.

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.

Twelve Changes a Producer Might Make

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.

What Each Mode Let Through

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.

A table headed ctr-demo.json and ctr-avro.json: verdicts against v1, titled twelve changes, three modes, one old consumer. One row per change, with BACKWARD, FORWARD and FULL verdicts and the records the old consumer failed to decode. Six changes pass all three. Add required field and humidity to float pass FORWARD only. Delete windspeed, hour int to long and add enum heavy_rain pass BACKWARD only, failing 17,379, 17,379 and 3 records. Rename temp passes none and fails 17,379. Beneath: Apache Avro 1.12.2 gave the same verdicts and counts.

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.

An isometric drawing of three blocks, heights to scale, headed ctr-report.json: heights to scale, out of 12 changes, titled what each mode let through. From left: BACKWARD, 9 of 12; FORWARD, 8 of 12; FULL, 6 of 12. Beneath: all 4 changes that moved the model passed every mode; BACKWARD also passed 3 of the 4 that broke the old consumer.

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.

Passing BACKWARD Does Not Protect Old Readers

Four changes left the old consumer unable to read some or all records. BACKWARD passed three of them. FORWARD and FULL passed none.

A flowchart headed BACKWARD passed these; the old consumer could not read them, titled passing BACKWARD does not protect old readers. The producer deploys v2; the consumer is still on v1, leads to delete windspeed, hour as long, or add heavy_rain. That leads to v1 finds no value and no default, cannot read a long as an int, or meets a symbol it lacks. That leads to the old consumer cannot decode the record. Beneath: failed records 17,379, 17,379 and 3; Confluent: under BACKWARD, upgrade all consumers first.

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.

Four Changes Decoded Cleanly and Moved the Model

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.

A dot chart headed ctr-demo.json: one dot per serving block of about 146 hours, titled four changes decoded cleanly and moved the error. Across 30 serving blocks from July to December 2012, the change in MAE in rides per hour, on a scale from -10 to 80, for four changes. Fahrenheit, humidity as percent and rename with defaults mostly sit between about 10 and 64, with a few blocks near or just below zero and Fahrenheit's last two near 75 and 80. Delete windspeed with defaults sits close to zero throughout. Beneath: means 30.20, 30.41, 36.39 and 0.27; clean MAE 44.5.

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.

A bar chart headed ctr-demo.json: the nine changes every serving record survived, titled share of predictions moved by more than one ride. Nine bars on a scale from 0 to 1, labelled optional, required, float, °F, percent, enum, season, wind 0 and rename 0. Only four bars show: °F near 1.0, percent near 0.98, wind 0 near 0.36 and rename 0 near 0.94. Beneath: °F 99.3%, percent 97.6%, rename with defaults 93.6%, windspeed with defaults 36.1%; float, one hour; the rest: none.

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.

A Default Turned a Loud Break into a Quiet One

The rename shows the trap best, because I ran it twice.

Two panels headed ctr-demo.json: the same rename, twice, titled a default turned a loud break into a quiet one. No defaults: rejected by all 3; had it shipped anyway, 17,379 of 17,379 records fail to decode. Defaults of 0: passed all 3; temp read as 0 every hour; MAE up 36.39 rides per hour. Beneath: worse in 27 of 30 blocks; from -7.2 to 63.6.

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.

A sidebar card headed after the results, titled humidity as a float. In large type: 1 of 4,376, predictions changed at all. Then: that one moved by 7.080 rides. Last: FORWARD passed it; v1 reads a float as a double.

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.

What a Check on Meaning Caught

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?

A hand-sketched column of four boxes joined by arrows, headed the semantic contract, written from the training rows before the run, titled four rules about meaning. 1, temp within 10 of its training range, -9.18 to 50.18. 2, feel_temp the same way, -10.0 to 60.0. 3, humidity from 0 to 1; windspeed from 0 to 10 above its training top. 4, every number field takes at least two values in a block. Beneath: clean blocks it fired on, 0 of 30; I wrote the rules knowing the twelve changes.

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.

A two-column table headed ctr-demo.json: the four changes that moved the model, titled the schema check passed them; the meaning check did not. Left, the change; right, semantic rule, blocks of 30. Temps in Fahrenheit, same name, same type: temp range 30, feel_temp range 30. Humidity as a percent, same name, same type: humidity range 30. Rename temp, with defaults of 0: temp constant 30. Delete windspeed, with defaults of 0: windspeed constant 30. Beneath, left: every mode passed all four. Right: each fired in every block.

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.

A line chart headed ctr-report.json: temps in Fahrenheit, after the results, titled even the cold blocks crossed the line. Across 30 serving blocks from July to December 2012, two lines on a scale from 30 to 110: the hottest hour in each block, falling from about 106 to about 57, and the coldest hour, falling from about 81 to about 42. A dashed line marks the rule's upper limit near 50. Beneath: limit 50.18; the coolest block's hottest hour read 57.1; many colder hours sat inside the 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.

Each Step Passed; the Whole Did Not

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.

A hand-sketched column of three boxes joined by arrows, headed ctr-demo.json: one change at a time, then all at once, titled each step passed BACKWARD; the whole did not. v1: windspeed is a number (double). v2: windspeed deleted; BACKWARD passes. v3: windspeed back, as text, default empty; BACKWARD against v2 passes. Beneath: v3 against v1 fails; a number cannot be read as text; only BACKWARD_TRANSITIVE checks this.

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.

A card headed ctr-report.json, after the results, titled a new value can wait months to break anything. In large type: 3 heavy rain hours in 17,379. Then: 2011, month 1, hour 16; 2012, month 1, hour 18; 2012, month 1, hour 1. Then: in the serving half, July to December 2012, 0. Last: the enum change failed only on those three records, months apart.

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.

Two Places to Check: the Merge and the Data

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.

Changes No Mode Allows

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.

A flowchart headed a change no mode allows, done on purpose, titled breaking changes get a new version and a plan. The change cannot be made compatible, leads to publish it as a new version, often on a new topic. That leads to write both old and new for a while. Then move each consumer, using the owner list. Then retire the old version when nobody reads it. Beneath: Confluent's docs: the more likely path is a brand-new topic, then moving applications to it.

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.

Why ML Teams Care Most

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.

What Teams Wrote About It

Data contracts became a popular idea around 2022. Here is what the sources I checked actually say, and no more.

A page of three entries, headed in each source's own words, checked 2026-10-01, titled what teams wrote about it. Convoy, 2022: contracts defined in the service's code and enforced in CI; the registry set to FORWARD, because producers update first; Convoy shut down in October 2023. Airbnb, 2020: data quality work that began with ownership that was not clearly defined. Google, Rules of Machine Learning: training-serving skew is a difference in performance between training and serving; one cause is handling data differently in the two pipelines.

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".

Four cards headed what each tool does, from its own documentation, titled where contracts are written and checked. Kafka, with its logo, with Confluent Schema Registry: Avro, Protobuf or JSON Schema; checks each new schema; default BACKWARD. dbt, with its logo, model contracts, 1.5 and later: checks column names and types before a model is built. JSON Schema, with its logo: required, enum, minimum and maximum, additionalProperties. Open Data Contract Standard, ODCS: a YAML format from Bitol, schema, quality, SLA, team, roles; v3.2.0. Beneath: ODCS, Confluent and Avro have no logo here; the logo set has none.

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.

Try It Yourself

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.

A real screenshot of VS Code with contract_demo.py open at the top of the file. The docstring asks which schema changes a registry lets through, and what each one does to a model that reads the data. It gives what the script needs, scikit-learn and pandas, the download of about 0.2 MB from OpenML, and how to run it. Then the design, written 2026-10-01 before the first run: the data, OpenML 42712, 17,379 hours of a bike share in 2011 and 2012; the producer and its v1 Avro schema, sending heavy rain as rain; the consumer, reading with v1 and feeding boosted trees trained on January 2011 to June 2012 and serving July to December 2012 in 30 blocks; the checker's rules from the Avro specification; the three modes; the twelve producer changes and the chain; and the start of what is measured for each change. The rest of the docstring and the code are further down.

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 Lab Report

A real terminal recording headed python contract_report.py, titled every number recomputed with my own code. It opens: bike share hours, OpenML 42712, 17,379, train 13,003, serve 4,376; every number recomputed with my own code, 444 checks agree; verdicts and decoded records match Apache Avro 1.12.2 for all twelve changes. Then four sections: what each mode let through, BACKWARD 9 of 12, FORWARD 8, FULL 6; a table of the old consumer's errors and MAE changes; the semantic contract, 0 clean blocks and 30 for each quiet change; and the facts found after the results.

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.

Five logo cards headed the tools, with their logos, titled what ran where. scikit-learn: the download and the boosted trees. NumPy: the blocks, the errors, the 32-bit float. pandas: reads the downloaded table for scikit-learn. Python: the checker and the decoder, no library. Apache Avro 1.12.2: a second opinion, in its own small environment.

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.

Check a Change Yourself

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 Code, Part by Part

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 .

How to Change a Schema, Step by Step

A hand-sketched column of six boxes joined by arrows, headed changing a schema, step by step, titled pick the mode, gate the merge, check the meaning. 1, write the contract down: shape, meaning, freshness, owner. 2, pick a mode that matches who upgrades first; consider _TRANSITIVE. 3, check every schema change in CI, against the registry. 4, be careful with defaults: a default can hide a missing field. 5, check the data's meaning at run time, not only its shape. 6, a change no mode allows: new version, both written, then retire. Beneath: here every mode passed a rename that fed the model 0 for temp.

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, or a Lighter One?

A two-column table headed grounded in this lesson's numbers, titled a full contract, or a lighter one? Left, worth it when: another team owns the data your model reads; many consumers read the same topic or table; a wrong value costs more than a stopped job; producers ship on their own schedule. Right, lighter is fine when: one team owns both ends and ships them together; the data is a one-off file for one analysis; a person reads every output before it is used; nothing downstream is trained or automated. Beneath, left: write it, and check it in CI and at run time. Right: a schema check and a test may be enough.

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.

What These Numbers Can and Cannot Tell You

A two-column page headed read before you trust these numbers, titled what these runs are, and what they are not. They are: one public dataset, bike share hours, 2011 and 2012; Avro's rules, my checker, and Avro's own library; one model, boosted trees; 30 blocks of serving hours; twelve changes I chose; semantic rules I wrote knowing the changes. They are not: not every table or every year; not Protobuf or JSON Schema; not every model; not a forecast of your errors; not every change a team makes; not proof such rules catch everything.

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.

What to Do Next

A hand-drawn list headed before the next schema change reaches your model, titled five questions for your own pipeline. Written?: is the contract written down, with units and an owner? Which mode?: which compatibility mode is set, and does it match who upgrades first? Defaults?: which fields have defaults, and what would a missing value become? Meaning?: what checks units and ranges at run time? Breaking?: what is the plan when a change no mode allows is needed? Beneath: here BACKWARD, the default, passed 9 of 12 changes.

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.

A closing card headed to keep, titled a shape check is not a meaning check. In large type: MAE +36.39 rides an hour. Beneath: a rename every mode passed, because both sides had defaults. Then: BACKWARD passed 9 of 12 changes; 3 broke the old consumer. Last: meaning rules, written knowing the changes, fired on all four quiet changes, in every block.

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.

Knowledge Check

Knowledge Check

4 questions - Score 80% to pass

Q1

Under BACKWARD, the registry's default, which reader does the mode make no promise to?

Q2

The producer started sending temperatures in Fahrenheit, with the same field names and types. What did the lab measure?

Q3

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?

Q4

Each step of the windspeed chain passed BACKWARD: v2 deleted it, v3 added it back as text. What does BACKWARD_TRANSITIVE add?

default

This is a real run in VS Code's terminal (python contract_demo.py).

A real screenshot of VS Code's terminal after running python contract_demo.py. It prints scikit-learn 1.9.1; hours 17,379, train 13,003, serve 4,376 in 30 blocks; heavy rain hours 3; clean MAE 44.5 rides an hour, 30.7 to 66.3; clean blocks a semantic rule fired on, 0. Then the verdicts, a row per change, with failed records of 17,379 for three changes and 3 for heavy_rain. Then what the model saw: Fahrenheit, MAE up 30.20; humidity as percent, 30.41; windspeed with defaults, 0.27; rename with defaults, 36.39; each with rules firing in 30 of 30 blocks. Last, the chain: v3 after v1 fails BACKWARD.

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.

decode
semantic
can_read