Data Engineering For Ml

Data Ingestion Pipelines: Getting Data Into the ML Platform

0 of 31 complete

0%

Contents

Back|Data Engineering For MlData Ingestion Pipelines: Getting Data Into the ML Platform
1/31
93 min left
  1. Home
  2. AI Engineering: Data, RAG and Agents
  3. Data Engineering for ML
  4. Data Ingestion Pipelines: Getting Data Into the ML Platform
Prerequisites
Batch vs Streaming Data for ML: When Fresh Beats Cheaprequired
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 31
Previous lessonBatch vs Streaming Data for ML: When Fresh Beats CheapNext lessonETL vs ELT: What a Model Loses When the Raw Rows Are Thrown Away

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 Clerk and the Daily Sheets

Imagine a small bakery chain with twelve shops. Every evening each shop sends the head office one sheet of paper: what it sold that day, hour by hour. A clerk at the head office copies every sheet into one big book. On Sunday the owner reads the book and decides how much bread each shop should bake next week.

Now four small things happen, and none of them makes a sound.

A courier hands over a sheet, but the clerk's receipt blows off the desk, so the next morning the courier comes back with the same sheet. The clerk copies it in again. That shop's Tuesday now counts twice.

One shop's sheet arrives on Monday, after the clerk has already closed last week's page. That day is simply missing from the book.

Another sheet arrives cut into strips and mixed up, because the shop printed it in pieces. The clerk copies the strips in the order they came out of the envelope. Next to each hour the clerk also writes "busier or quieter than the hour before", and now that note compares each hour with a random other hour.

A fourth shop starts using a new form, where the box that used to say "loaves" now says "bread, units". The clerk looks for "loaves", does not find it, and writes nothing.

An illustration of a woman sitting in an armchair, typing on a laptop, next to text. Headed the clerk and the daily sheets, titled four quiet ways a copy goes wrong. Beside her: every day, a file of the day's records; someone copies each one into the table a model learns from. Below, four rows. Sent twice: 66 of 756 daily files landed twice, 3,168 rows counted twice. Late: 66 files missed the table, 3,168 rows gone. Shuffled: 66 files landed out of order, 1,406 lag inputs wrong. Renamed: a column renamed upstream, 3,609 rows with that column empty. Last: in the lab, none of the four raised an error; the table just came out wrong.

Every one of these leaves the book looking normal: neat handwriting, the right pages, numbers in every column but one. The owner will bake too much or too little, and nobody will know why. This lesson is about that clerk's job in a data platform, where the sheets are files of records and the book is the table a model learns from. The four numbers in the picture come from the lab later in this lesson. There I made each of these mistakes happen on purpose, and measured what it did.

Where This Lesson Starts

The previous lesson compared two ways to move data: in batches on a schedule, and as a stream, one event at a time. This lesson looks at the step that comes first in both: copying data out of the systems that run a business and into the place where models learn. That step is called ingestion.

In most companies, the data a model needs lives in systems built for something else. Think of the shop's order database, the app that records every tap, or a partner that sends a file every night. Ingestion is the work of copying that data out, safely and completely, into storage the data team controls. Everything a model learns later is built from that copy, so a mistake at this step reaches every number built from it and every model trained on it.

The lesson has two halves. The first half explains the ideas: where data comes from, who starts each transfer, and the three common ways to copy it. It also covers the problems that make ingestion hard: copies, bursts, late data and changing columns. The first version of this lesson stopped there. For this version I checked every claim in it against the tools' own documentation, and fixed or cut the ones that did not hold. The list of claims, sources and verdicts is in the lab folder, in scripts/labs/dataeng/results/ip-factcheck.json.

The second half measures. I took real data, sent it through a small copying pipeline one file a day, and broke the pipeline in four quiet ways, the clerk's four mistakes. Then I asked three questions. What did each fault do to the table? What did it do to a model trained on that table? Which simple check would have noticed?

The answer surprised me. The model's accuracy was a poor alarm. It went up for some faults and down for others. Retraining the same model on the same clean table, with only the random seed changed, moved it more than most faults did. Six cheap checks on the table itself, which need no right answers and no model, caught each of the four faults in every run. I wrote each check knowing which fault it was for, so that part was expected; the limits slide says what they would miss.

Ten Words for This Lesson

A hand-drawn list headed ten words for this lesson, titled the pieces between a source and a model. Ingestion: copying data from the systems that run a business into the place models learn from. Source: the system the data comes from: a database, an app, a partner's files. File: here, one day of records, sent as one piece. Landing zone: storage that keeps each file exactly as it arrived, also called bronze. Key: a value that names one record and no other: here, the half-hour's row number. Duplicate: the same record, landed more than once. Idempotent: a write that gives the same result whether it runs once or twice. Event time: when a thing happened; processing time is when it reached us. Schema: the names and types of the columns a file should have. Lag input: the label of the half-hour before, given to the model as an input. Beneath: none of the lab's four faults raised an error; the table just came out wrong.

Ingestion is copying data from the systems that run a business into the place models learn from. The system it comes from is the source. In this lesson's lab the source sends one file a day: that day's records, as one piece. The first place a file lands is the landing zone: storage that keeps each file exactly as it arrived. Many teams call it the bronze layer.

A key is a value that names one record and no other, like an order number. When the same record lands more than once, the extra one is a duplicate. A write is when doing it twice gives the same result as doing it once. Setting a box to 5 is idempotent; adding 5 to it is not.

Event time is when something happened; processing time is when it reached us. A schema is the list of column names and types a file should have.

The lab's model predicts, for each half-hour, whether an electricity price went up or down. That right answer is the label. The model is also given the label of the half-hour before as one of its inputs, which I call the lag input. A feature is any input a model reads. Accuracy is the share of rows a model got right, and a point is one hundredth of accuracy. Here is the path those words describe, as the lab builds it.

Why the Model Never Reads the Live Database

Here is a story I made up to show the problem. I have not measured it, but each part of it is how these systems behave.

A data scientist builds a recommendation model. It needs three inputs: how many orders each user placed this week, their average basket size, and whether they opened the app today. All three live in the production Postgres database that runs the shop. So the first version of the model simply queries that database each time it makes a prediction.

It works in the demo. Then it goes live, and checkout slows down. The model's big analytical queries compete with real customers for the same processors and the same disks. A long-running read also holds back Postgres's clean-up of old row versions, a job called VACUUM, so busy tables grow bloated.

One thing does not happen: in Postgres a plain read never blocks an ordinary write, an INSERT, UPDATE or DELETE. Its documentation says "reading never blocks writing and writing never blocks reading", because readers and writers see separate versions of each row. There is one exception worth knowing. Most forms of ALTER TABLE need the strongest lock, and a running read does block that. While the ALTER TABLE waits, the writes that arrive after it wait behind it. The database team bans the model from production anyway, and the project stalls.

So a model should not read the live product database, and it should not depend on a source it does not control. Instead, you copy the data you need into the platform, on your own schedule, in your own shape, without disturbing the systems that run the business. That copying is ingestion, and every later lesson in this chapter builds on it.

There is a second, quieter reason ingestion matters. It sets the freshness ceiling of every feature: the freshest a value can ever be. A feature computed from a source you copy once a night can never be fresher than a night old, however fast everything above it runs. That ceiling is set at the very first copy, and nothing downstream can raise it.

Four Sources, One Landing Zone

Every piece of data a platform ingests comes from one of four kinds of source, and each behaves differently.

  • Databases (Postgres, MySQL, MongoDB). Structured, changing all day, and not for heavy reads in production.
  • APIs (a payments provider, a weather service, a partner catalogue). An API is a door another company's system opens to programs: you ask, it answers, one page at a time, within a .
  • Events (clicks, purchases, ride requests). Behaviour the app records as it happens, which often never lives in a table.
  • Files (daily CSV or Parquet exports, vendor feeds, log dumps). Dropped into storage by systems you do not control. The lab's source is this kind.

Whatever the source, the first stop is the same: the landing zone. This is cheap object storage, usually Amazon S3 or something like it, where data lands exactly as it arrived, untouched. Keep two times with every record: when it happened, and when it arrived.

An editorial page in five labelled zones, headed four kinds of source, one first stop, titled land it as it came, then clean it. Databases, with the PostgreSQL, MySQL and MongoDB logos: rows that change all day; heavy reads compete with the business, so copy the changes out. Events, with the Apache Kafka logo: clicks and orders the app publishes as they happen, often never in a table. APIs: an outside service you can only ask, a page at a time, within a rate limit. Files, with the Amazon S3 logo: daily exports dropped by systems you do not control; the lab's source is this kind. Landing zone first: every file exactly as it arrived, with when each record happened and when it arrived; clean, remove copies and build features only after. Beneath: the untouched copy is what lets you fix a bug and rebuild, without asking the source again.

The rule that keeps this sane: land raw first, transform later. Storage is cheap per gigabyte, and that untouched copy is what makes recovery possible. If a cleaning job has a bug, you fix the job and rebuild from the landing zone. You do not go back to the source, because the source may have changed.

An order a customer edited yesterday looks different today, and a click a phone sent while offline is long gone from the app's memory. The landing zone is your one lasting record of what actually arrived. In the lab, every fix at the end works only because each landed file was still there to rebuild from.

Pull or Push: Who Starts the Transfer

Before you choose between batch, CDC and streaming, settle a simpler question: who starts the transfer? There are only two answers, and the choice decides who controls the timing and who carries the load.

A two-column page headed who starts the transfer, titled pull on a clock, or push on each change. Left, pull: the platform asks the source what changed, on its own schedule; fresh only up to the last run; a read on the source every run; tracks a watermark such as updated_at, and a deleted row has none, so it is missed. Right, push: the source sends each change as it happens, into a broker; changes arrive in seconds; little load: an event, or a reader of the database's log; the consumer still pulls from the broker, and tracks its offset. Beneath, left: fits outside APIs and slow tables. Beneath, right: fits busy databases and events.

In a pull, the ingestion job reaches into the source and asks what changed, on a clock the platform owns. Batch jobs and API polling both work this way. The platform controls the timing, which is convenient. But features are only as fresh as the last run, and every run puts a read on the source that competes with real traffic.

To avoid reading everything each time, a pull keeps a watermark, usually a column like updated_at. Each run fetches only the rows changed after the highest value it saw last time. A deleted row has no updated_at left to find, so a watermark pull never sees deletes. Debezium's own article on this says polling "will not allow you to identify any records that have been deleted since the last poll".

In a push, the source sends each change as it happens, and the platform only has to receive it. Event streaming and change data capture both work this way. Both usually deliver into a broker: a system such as Apache that stores messages until readers take them. Changes flow in seconds, and the source carries little extra load.

A few Kafka words first. A topic is a named stream of messages inside the broker. A producer is a program that writes messages to a topic, and a consumer is a program that reads them.

One detail is easy to miss. Even in a push design, the consumer pulls from the broker. Kafka's documentation says data is "pushed to the broker from the producer and pulled from the broker by the consumer". That lets a slow reader fall behind without breaking anything. The reader keeps its place with an : the position of the last message it finished. After a crash it carries on from there, instead of starting over or skipping ahead.

The Connector Follows the Source

You rarely write ingestion by hand. You pick a connector, a ready-made program that knows how to read one kind of source, and the one you pick follows from the kind of source. Managed tools such as Airbyte and Fivetran each list hundreds of connectors. Debezium, the open-source tool for reading database change logs, covers about a dozen databases, in depth.

Six brand cards headed the connector follows the source, titled real tools, one per kind of source. PostgreSQL, with its logo: a log reader such as Debezium turns each change into an event. MySQL: the same idea, reading MySQL's own change log, the binlog. MongoDB: change streams, which MongoDB builds on its oplog. Apache Kafka: a consumer reads each partition and commits its offset. Amazon S3: a watcher picks up each new file and lands it as it is. Airbyte: ready connectors for APIs and databases: pages, retries, rate limits. Beneath: Debezium has no logo in the icon sets I use, so it is named, not drawn.

The reason to use a connector is not the easy part. Reading rows from a database or fetching a page from an API is simple. The value is in the edges. A good connector works through a large API result a page at a time. It waits and retries when an API answers 429, the code for "too many requests". It retries a dropped connection without losing its place, notices when the source's columns change, and saves its progress so a restart carries on cleanly.

Airbyte's documentation gives one example, for connectors built with its Connector Builder. When no error handling is set, they retry a request that got a 429 or a server error five times, waiting longer each time. Those are exactly the details a script written in an afternoon forgets, until the night it breaks. Write the connector yourself, and you volunteer to learn each of them again.

Three Ways to Extract: Batch, CDC and Streaming

Once you know the source and who starts the transfer, you pick how to extract. There are three common ways, and the choice quietly sets the freshness ceiling from two slides back.

A three-row table headed three ways to extract, titled the mode sets the freshest a feature can be. Batch pull: on a schedule, ask for rows changed since the last run; fresh up to the schedule; a watermark misses deleted rows. Change data capture: read the database's own change log; fresh in seconds; sees deletes; an abandoned reader keeps the log on disk. Event streaming: the app publishes each event as it happens; fresh in seconds or less; needs a broker to run and watch. Beneath: no layer above ingestion can make a feature fresher than the first copy.

Batch pull is the simplest. On a schedule, you ask the source for everything changed since last time, usually by tracking a watermark like updated_at. Hourly, nightly, whatever the clock says. It is easy to reason about and cheap to run, but features are only as fresh as the schedule, and a watermark misses deleted rows.

Change data capture (CDC) reads the database's own change log and turns every insert, update and delete into an event. In Postgres that log is the write-ahead log, or WAL; in MySQL it is the binlog. For MongoDB, Debezium reads change streams, which MongoDB builds on its own log, the oplog. CDC reads the log instead of querying tables over and over. So it puts little load on the source, it sees deletes, and changes arrive close to real time.

Event streaming captures things as the application does them, publishing to a broker such as the moment a click or a purchase happens. It is the usual way to ingest behaviour that never touched a database; log files shipped to storage are another.

ModeFreshest possibleLoad on the sourceSees deletes

How CDC Actually Works

CDC feels like magic the first time you see it, so here it is slowly. The surprising part is how little the source has to do.

A sequence diagram with four columns: the app, Postgres, Debezium, Kafka. Headed change data capture, one change, titled the database writes its log anyway; a reader turns it into events. Step 1, the app sends Postgres an ordinary UPDATE. Step 2, Postgres, to itself: the change goes into its log. Step 3, Postgres sends Debezium the committed changes, by the slot. Step 4, Debezium sends Kafka an event: op, before, after. Step 5, Debezium tells Postgres it has read up to here, so the log can go. Beneath: a consumer then pulls events from Kafka and commits its offset; a crash before the commit means those events come again.

The application writes to Postgres exactly as it always has. It has no idea CDC exists. Postgres, for its own crash recovery, writes every change into its write-ahead log first, including changes from transactions that later roll back. A tool like Debezium holds a slot: a named bookmark in that log that Postgres keeps for one reader.

Through the slot, Postgres hands over each transaction once it commits. By default, its documentation says, changes are "only passed to the output plugin at commit (and discarded if the transaction aborts)". Since Postgres 14, a plug-in can also choose to receive a large transaction in pieces before it commits; that is an option, not the default. After a first snapshot, which does read each table once with an ordinary SELECT, Debezium reads only the log to find changes.

Each change becomes an event that says what happened (op: create, update or delete) and carries the row's values. after is the new row. before holds the old row in full only when the table's REPLICA IDENTITY setting is FULL. With the default setting, update and delete events carry old values for the primary key columns only. PostgreSQL's description of its replication messages adds that for an update, even the old key is sent only when the update changed it. That matters if a feature needs to know what a value used to be.

A note on sources. Debezium's documentation describes REPLICA IDENTITY the same way, but its worked example of an update shows the key in although the key did not change. The wording above follows PostgreSQL's own pages on ALTER TABLE and on the logical replication message formats. Check what your version sends before a feature relies on .

What CDC Looks Like in Config

You rarely write CDC by hand. You configure a connector and let it run. Debezium is the open-source engine most teams reach for. Managed platforms like Fivetran and Airbyte wrap the same idea, if you would rather not run the plumbing yourself. Here is a trimmed Debezium connector config that captures changes from two Postgres tables and publishes them to .

{
  "name": "orders-cdc-connector",
  "config": {
    "connector.class": "io.debezium.connector.postgresql.PostgresConnector",
    "database.hostname": "orders-db.internal",
    "database.port": "5432",
    "database.user": "debezium",
    "database.dbname": "shop",

    "plugin.name": "pgoutput",
    "slot.name": "debezium_orders",
    "publication.autocreate.mode": "filtered",

    "table.include.list": "public.orders,public.order_items",
    "topic.prefix": "shop",

    "snapshot.mode": "initial",
    "tombstones.on.delete": "true",
    "decimal.handling.mode": "string"
  }
}

A few lines carry most of the weight. plugin.name chooses how Postgres turns its log into messages. Debezium's default is decoderbufs, a plug-in you must install on the server. pgoutput is built into Postgres 10 and later. The documentation of Amazon RDS, Azure Database for PostgreSQL and Google Cloud names other plug-ins, such as wal2json, but not decoderbufs. So pgoutput is the one of Debezium's two choices that all three offer, and it is set here on purpose. slot.name is the slot that remembers how far Debezium has read, so a restart carries on cleanly.

snapshot.mode: initial matters on the first start, when the connector has no stored position. It reads every row of each table as it is at that moment, then switches to reading the log from that point. You get the current state plus every change after it, without a gap, but not the past edits of each row.

When Data Outruns You: Throughput and Backpressure

CDC and streaming solve freshness, but they bring a problem batch never had. A source can produce faster than the pipeline can write. During a sale or a dinner rush the buffer between them fills, and what happens next decides whether ingestion survives the burst.

A hand-sketched flow headed when events come faster than they are written, titled a durable buffer turns a burst into a delay. A box, producer: a burst, points to a cylinder, broker: the backlog, kept on disk, which points to a box, consumer: its own pace, which points down to a cylinder, landing zone. Below, three levers. Keep it: the broker holds the backlog, but only for its retention time. Add readers: more partitions and consumers; each partition goes to one consumer in a group. Shed: drop or sample low-value events on purpose, and count what was dropped. Beneath: without one of these, a burst ends in a crash or in events lost with no record.

is the signal that runs the other way, from the overwhelmed reader back toward the writer. When the buffer passes a high-water mark, the system does one of two things. It tells the writer to slow down, or it keeps the backlog in lasting storage and lets the reader catch up at its own pace. Either way the burst is absorbed instead of dropped. Without it, a spike ends in one of two ugly ways. The reader runs out of memory and crashes. Or events are thrown away with no record that they ever existed, which quietly skews every feature built from that window.

There are three levers, and mature platforms use all three.

Keep it. A broker such as holds the backlog on disk, so a slow reader only lags, but only for as long as the topic keeps data. Kafka's default retention is 168 hours, a week. A reader further behind than that has lost those events, and with the default setting it jumps ahead to the newest data without an error.

Add readers. A topic is split into partitions: separate, ordered pieces of its log that can be read in parallel. A consumer group is a set of consumers that share the reading of one topic. More partitions and more consumers raise , with one limit: each partition goes to exactly one consumer in a group, so extra consumers get none.

Shed. Under extreme load, when freshness matters more than completeness, drop or sample low-value events on purpose. Count what you dropped, so the loss is on record. A durable buffer plus a slow-down signal turns a possible outage into a delay that drains on its own once the surge passes.

Duplicates Are Not Optional, So Plan for Them

Any pipeline that retries can deliver a record more than once, and you cannot know in advance which records. 's default is at-least-once. A producer retries after a timeout, sometimes when its first try already worked but the acknowledgement never came back. In current Kafka the producer is idempotent by default, so such a retry does not write a second copy. Its settings say "Idempotence is enabled by default if no conflicting configurations are set". When idempotence is off, or turned off by conflicting settings, retries "may write duplicates of the retried message".

Idempotence in the producer does not help further along the pipe. A connector that restarts, or a consumer that reads again after a crash, can still deliver a record twice. And a partner re-sends yesterday's file because their job failed halfway.

If your pipeline simply appends, the training table quietly counts some things twice. A feature like "orders this week" then reads high for exactly the users whose records were retried. The pipeline never errors, and the number looks believable. Here is what that looks like in the lab, where the connector sent a tenth of the daily files twice and the job appended both copies.

A hand-drawn sketch headed sketched: the first four daily files of 1996, at 10% and seed 0, titled a file sent twice lands twice. Three rows of four boxes. File for: 7 May, 8 May, 9 May, 10 May. Sent: 48, 48, 48, 48. Landed: 48, 48, 96, 96, the last two marked. Below: 9 May and 10 May were touched: each file came twice. The count check: 96 landed, 48 sent. Beneath the sketch: over all 756 files, 66 landed twice, 39,417 rows for 36,249 records, 3,168 of them copies; no step failed.

On the first four days of the data, two files came twice. The only visible sign was a count: 96 rows landed where the source had sent 48. Over the whole training window, 66 of the 756 files landed twice, so the table held 39,417 rows for 36,249 real records, and 3,168 of them were copies. No step failed. (The calendar here is rebuilt from row numbers, 48 half-hours a day from 7 May 1996, as the lifecycle chapter does.)

The defence is : make writing the same record twice give the same result as writing it once. You do that by giving every record a stable key, and by turning every write into a keyed upsert instead of a blind insert. An upsert updates the row if its key already exists, and inserts it if not.

The key is usually the source's primary key, like order_id. A version then decides which copy wins. Putting the version into the key would store every version as a separate row, which brings the duplicate problem back. The first time you see a key, you insert. When you see it again, you keep this copy only if its version is newer, because retries and out-of-order delivery mean a copy may be newer or older. Without that check, a late old copy would overwrite a newer value and quietly roll a row back in time.

Late Data and Backfills

Data does not always arrive when it happened. A phone is offline for two hours, then syncs. A partner batches its exports and sends them at midnight. A queue backs up during a traffic spike. This gap between event time, when the thing happened, and processing time, when you received it, is one of the trickiest problems in ingestion. The two terms come from Tyler Akidau's essays on stream processing, "Streaming 101" and "Streaming 102", published by O'Reilly.

Watch what goes wrong if you ignore it. An hourly job computes "features for the 10:00 window" at 11:00. It looks complete. Then at 12:30 an offline phone syncs its 10:00 click. If the window is built from what had arrived by 11:00, and never built again, the 10:00 features stay wrong for good.

If you train on that window, you teach the model from incomplete data. The error is invisible, because the window looked finished when you sealed it. The lab's late files are the same thing at the size of a day. There, 66 files had not arrived when the training table was built, so 3,168 rows were simply missing.

Two rules fix it. First, record both times with every record, and compute windows by event time, so a late record still counts in the window it belongs to. Many teams store the landing zone by arrival time, because it only ever grows at the end and device clocks can be wrong. They then group by event time in the next layer. Second, treat any feature computed near "now" as provisional, and run a backfill once late data has settled: re-run past windows to correct them.

In practice you also set a bound, a watermark in the streaming sense. This is a different use of the same word as the pull's watermark. Akidau defines it as "a notion of input completeness with respect to event times". In plain words: we will wait this long for stragglers, then finalise. You cannot keep every window open forever for a phone that may never reconnect.

A backfill is worth seeing as its own shape of run. The normal pipeline is incremental: each cycle takes one small, roughly constant slice of new data. A backfill is the opposite: one large sweep back over history. You run it to correct windows that late data changed, or to fill a new feature across all of its past.

It needs care an incremental run does not. A sweep over months of the landing zone can flood a broker, overload a database or run up a cloud bill if you run it at full speed. So limit its rate and run it off-peak. And lean on the from the last slide. Because the writes are keyed, a backfill that overlaps the live run lands in the same correct state, instead of fighting it.

Schema Drift: When the Source Changes Underneath You

Your ingestion runs perfectly for months. Then one morning it either crashes or, far worse, quietly starts writing empty values. A backend team shipped a change. They renamed a column, changed a type or dropped a field, and nobody told the data team, because nobody knew the data team depended on it. This is schema drift. Its cause is mostly organisational. The team that owns the source and the team that uses it are usually different, and the source team ships on its own schedule. Where the damage lands depends on a choice you made long before: where you enforce the shape of your data.

A two-column page headed where the shape is enforced, titled check it on the way in, or when it is read. Left, schema on write: check each file's columns before storing it; bad data never lands; a rename stops the job, loudly. Right, schema on read: store each file as it came, and apply the shape when it is read; nothing stops a file from landing; a rename lands as an empty column, quietly. Beneath, left: use it on the way into the trusted table. Beneath, right: use it for the landing zone.

Schema on write checks each file at the door, so a breaking change is caught early, but it stops the pipeline. Schema on read lands everything and applies the shape only when the data is read. Ingestion never stops, but bad data can sit unnoticed in the landing zone. The sturdy answer is to use both, at different layers. Use schema on read into the landing zone, so nothing can stop a raw file from landing. Use schema on write into the trusted table, so no unchecked row is used downstream.

The lab's rename shows why the second half matters. Its job read columns by name with pandas, and pandas does not complain about a missing name. reindex fills it with an empty value, NaN (pandas' documentation: values with no match "are assigned NaN"), and the job then filled that with 0. From the first renamed file on, 3,609 rows had a demand of 0, and nothing failed.

Change at the sourceSafe?What ingestion should do

Bronze, Silver, Gold, and Real Freshness

Raw ingested data is honest but messy, so you refine it in layers. Databricks calls this the medallion architecture. In its documentation, bronze "contains raw, unvalidated data". Silver holds "validated, cleaned, and enriched versions of the data". Gold holds "highly refined views of the data" that drive dashboards and machine learning.

In this lesson's terms, bronze is the landing zone. Silver is bronze after removing copies with keyed upserts and checking the schema: the first layer worth trusting. Gold is the features, computed with versioned logic, so a feature means the same thing everywhere it is read. Versioning that logic matters, because the day someone quietly redefines "active user" is the day training and serving stop agreeing.

From gold you publish to a with two copies, fed from the same place. The offline store is for training, and is read over months of history in bulk. The online store is for serving, and returns one user's features in milliseconds. Feeding both from one place prevents train/serve skew, where a model learns from one definition and predicts with another. The course measures it in train/serve skew.

Now the point that ties the first half together. Freshness is decided at ingestion. A nightly batch pull can never give you gold features that are seconds old. However fast the layers above it run, the ceiling was set at the first copy.

Three companies have written about this in their own engineering blogs.

DoorDash built a framework it calls Riviera on Apache Flink, a stream-processing engine. With it, teams describe real-time features in a configuration file. It reads from Kafka and writes to stores including , which powers DoorDash's feature store.

Uber's Michelangelo platform computes some features in batch and some near real time. Its example comes from Uber Eats: a restaurant's average meal preparation time over the last seven days comes from a batch job. The same average over the last one hour comes from a streaming job that reads .

Netflix carries its event data through Kafka, in its Keystone pipeline. Separately, it wrote that it builds training features from dated snapshots of its online services. That way a model learns from the inputs it would have seen at that moment.

The pattern is the same each time. Pick the ingestion mode that matches how fast the decision must be, and land the data raw. Remove copies with keyed upserts, and absorb bursts with a durable buffer. Correct late windows with backfills, and check the shape at the door of the trusted layer. Only then let it become a feature. The second half of this lesson measures what happens when a pipeline skips some of that.

What the Lab Ran

I wrote the lab's design at the top of its file, scripts/labs/dataeng/ingest_lab.py, before it ran. Everything below comes from its stored results, results/ingest.json, and from a follow-up I designed after seeing them, which the slides mark as such.

An editorial page in four labelled zones, headed what the lab ran: ingest_lab.py, designed before it ran, titled one source, one job, four faults. The source: Elec2, 45,312 half-hours, sent as one file a day; 756 files hold the 36,249 training rows, and the last 9,063 rows are the clean test rows. The job: selects the columns by name, fills anything missing with 0, and takes the lag input from the row before it in the table. The model: boosted trees with the lag input, 0.8193 on a clean table; without the lag, 0.7527. The faults: each touches 5%, 10%, 25% or 50% of the files, with 5 seeds for which files; the headline is 10%, seed 0. Beneath: the follow-up, with early stopping off, was designed after the main results.

The data is Elec2: 45,312 half-hours of the New South Wales electricity market, from 7 May 1996 to 5 December 1998, 48 a day, in time order. It is small, public and real, and scikit-learn downloads it from OpenML. For each half-hour the label is whether the price went UP or DOWN against its average over the last 24 hours. The first 36,249 half-hours are the training window, and the last 9,063 are the test rows. This is the same split as the lifecycle chapter's lesson on stale pieces in a pipeline, so its numbers can be checked against this one.

The source and the job are a simulation, and I say so plainly. The source sends one file a day: 756 files cover the training window, the last one short. The job is written the way many first versions are. It reads the files in the order they landed, and picks the eight input columns by name. It fills anything missing with 0. It makes the lag input by taking the label of the row above it in the table. In a clean table the row above is the half-hour before, so the lag is right.

The model is a set of boosted trees: many small trees of yes-or-no questions, built one after another, each one correcting the mistakes of the trees before it. It is scikit-learn's HistGradientBoostingClassifier, given the eight inputs plus the lag. On a clean table it scored 0.8193, and the same trees without the lag input scored 0.7527. The first number matches the lifecycle chapter's stored result exactly, and the lab refuses to save anything if it does not.

What Each Fault Did to the Table

Here is the training table each fault built, read from results/ingest.json, with each fault touching 10% of the files and seed 0.

A table headed ingest.json: the training table each fault built, 10% of files, seed 0, titled two faults changed the row count; two did not. Columns: files, rows, copies, missing, lag off, no demand. Clean: 0, 36,249, 0, 0, 0, 0. Sent twice: 66, 39,417, 3,168, 0, 27, 0. Late files: 66, 33,081, 0, 3,168, 27, 0. Rows shuffled: 66, 36,249, 0, 0, 1,406, 0. Column renamed: 76, 36,249, 0, 0, 0, 3,609. Cells that differ from clean are outlined. Beneath: counts of rows; lag off: lag inputs that are not the true previous label; no demand: rows whose nswdemand was missing and became 0.

Sent twice landed 39,417 rows for 36,249 records: 66 files came twice, and 3,168 rows were copies. Late files landed 33,081 rows: 66 files had not arrived, and 3,168 records were missing. Both also got 27 lag inputs wrong, at the edges of files. After a copy or a gap, the "row above" a file's first row is not the half-hour before. Rows shuffled landed every row once, but 1,406 lag inputs were wrong. Column renamed landed every row once with every lag right, and 3,609 rows had no demand value, which became 0. The most important number of this slide is not in that table, though.

Two panels headed ingest.json: all 64 fault runs, every share and seed, titled nothing failed; something always noticed. A step raised an error: 0 of 64 runs; every table was built and trained on. A table check fired: 64 of 64 runs; at least one of six checks, with no labels. Beneath: each check was written for one of these faults, so this was expected; it does not mean they catch every fault.

Across all 64 fault runs, every share and every seed, no step raised an error. Every table was built, and every model trained on it; the job did what it was written to do. And in every one of those 64 runs, at least one of six simple checks on the landed files fired. That was expected, since I wrote each check for one of these faults. None of those checks needs a label or a model. It is tempting to think a row count is enough, so here are the row counts on their own.

An isometric drawing of five blocks, heights to scale, headed ingest.json: rows in the training table, 10% of files, seed 0, titled a row count sees two of the four faults. From left to right: 36,249, clean; 39,417, sent twice, the tallest; 33,081, late files, the shortest; 36,249, rows shuffled; 36,249, column renamed. Beneath: shuffled and renamed both land exactly 36,249 rows, so only an order check and a schema check can see them.

Why a Shuffle Breaks the Lag Input

The shuffle is worth a closer look, because it breaks something the other faults barely touch. keeps messages in order only within one partition. Its introduction says it "only provides a total order over records within a partition, not between different partitions in a topic". A day's rows spread over several partitions can come back interleaved, though each partition keeps its own rows in order. My simulation puts the rows of each touched file in a fully random order, which is the worst case of that.

A hand-drawn sketch headed sketched: the first shuffled file, for 9 May 1996: its first 8 rows as they landed, titled the row before is not the half-hour before. Four rows of eight boxes. Half-hour: 1130, 0600, 2100, 0800, 1730, 2130, 2300, 0630. Label: U, U, U, D, U, U, U, U. Lag used: U, U, U, U, D, U, U, U, the third, fourth and fifth marked. True lag: U, U, D, D, U, U, U, U. Below: U means UP, D means DOWN; 1130 means 11:30; lag used is the label of the row above, as landed. Beneath the sketch: here 3 of 8 lag inputs are wrong; the whole day, 19 of 48; the whole table, 1,406 of 36,249; the first half-hour of a day is counted as 00:00.

This is the first shuffled file, the one for 9 May 1996. The label of each row is right, and nothing about the rows themselves changed. What changed is the row above each one. The job's lag input is "the label of the row above", which is the half-hour before only if the rows are in time order. Here, 3 of the first 8 lag inputs are wrong; over the day, 19 of 48; over the table, 1,406 of 36,249. As in the lifecycle chapter, the data does not say whether a day's first half-hour starts at midnight, so I count it as 00:00.

The lag input matters because the model leans on it heavily: the trees scored 0.8193 with it and 0.7527 without it. Here is one possible reason the shuffle hurts, which is my reading and not a separate measurement. Taught from a table where the lag is sometimes wrong, the trees learn to trust it a little less. They then pay for that caution on test rows, where the lag is always right. The fix is simple once you see it: build any input that depends on order in key order, never in landing order.

Which Check Caught Each Fault

The six checks were fixed in the design, before the run, and each one reads only the landed files. Row count compares each file's rows with the count the source says it sent, its manifest. Copies looks for any key that appears twice. Missing keys looks for any key of the training window that never arrived. Key order looks, in landing order, for a key that is not larger than the key before it. Column names compares each file's columns with the expected list. Empty column looks for an input with no value before the job fills it.

A table headed ingest.json: runs in which each check fired, all shares and seeds, titled each fault set off a different check. Rows: row count, copies, missing keys, key order, column names, empty column. Columns: clean, twice, late, shuffled, renamed. Row count: no, 20/20, 20/20, 0/20, 0/4. Copies: no, 20/20, 0/20, 0/20, 0/4. Missing keys: no, 0/20, 20/20, 0/20, 0/4. Key order: no, 20/20, 0/20, 20/20, 0/4. Column names: no, 0/20, 0/20, 0/20, 4/4. Empty column: no, 0/20, 0/20, 0/20, 4/4. Cells that fired are outlined. Beneath: each check was written for one of these faults; each that fired for a fault fired in every one of its runs, and none fired on the one clean table.

Each fault set off a different check, and each check that fired for a fault fired in every one of that fault's runs, at every share and seed. None fired on the clean table. Sent twice set off the row count, the copy check and key order. Late files set off the row count and missing keys. Rows shuffled set off only key order. Column renamed set off column names and the empty-column check.

Two things follow. First, no single check is enough: the row count, the check teams add first, misses half of these faults. Second, all six together are cheap. Each is a count or a comparison over keys and column names that a pipeline already has in hand. In the lab they need no labels, no model and no waiting for the right answers. So they can run on every file, the moment it lands.

One caution before you trust that table. I wrote each check together with the fault it matches, so each one firing in every run of its fault was expected, not a discovery. The limits slide lists what these six would miss, such as a change of units.

What the Model's Accuracy Said

Now the model. For each run the lab trained the same trees on the faulty table and scored them on the 9,063 clean test rows. At the headline, 10% of files and seed 0, the clean trees scored 0.8193. The trees scored 0.8230 with the copies and 0.7849 with the late files. They scored 0.8143 with the shuffle and 0.8134 with the rename. The copies raised accuracy. Before reading anything into that, look at the model on its own.

Two panels headed ingest.json: the clean table, no fault at all, titled the same table, another seed, another accuracy. Random_state 0: 0.8193, the lab's clean model. Random_state 1 to 9: 0.7884 to 0.8322, the same rows, the same trees, another seed. Beneath: with early stopping on, which it is by default here, the trees set aside a random 10% of the rows, so the seed changes what they learn from; they never stopped early, using all 100 rounds every time.

I trained the same trees on the same clean table nine more times, changing only random_state, the seed of the model's own random choices. They scored from 0.7884 to 0.8322, a spread of 4.38 points with no fault at all.

The reason is a default I had not thought about when I wrote the design.

With more than 10,000 rows, scikit-learn's boosted trees turn on early stopping. They set aside a random 10% of the training rows to decide when to stop adding trees, and never learn from those rows. The seed chooses which rows are set aside. So does any fault that changes the number of rows or the order of the labels, because the split is drawn from the row count and the labels. Early stopping itself never stopped early: every model used all 100 of its rounds, so its only effect here was the 10% it set aside.

A dot chart headed ingest.json: accuracy on the 9,063 clean test rows, every fault run, titled most fault runs landed inside the clean model's own spread. Four groups of dots on a scale from 0.76 to 0.86, with dashed lines for clean near 0.82, clean best seed near 0.83 and clean worst seed near 0.79. Twice: 20 dots from about 0.79 to 0.84. Late: 20 dots from about 0.78 to 0.83. Shuffled: 20 dots from about 0.78 to 0.81, none above the clean line. Renamed: 4 dots near 0.81. Beneath: twice 0.7865 to 0.8387; late 0.7849 to 0.8347; shuffled 0.7811 to 0.8143; renamed 0.8083 to 0.8165; 51 of 64 runs sit inside the clean table's seeds 1 to 9, 0.7884 to 0.8322.

Against that spread, most fault runs look like noise: 51 of the 64 landed inside the clean model's own range. The copies scored above clean in 9 of 20 runs and below in 11. Late files scored below clean in 17 of 20. The shuffle and the rename never scored above clean (the rename had only 4 runs).

Holding the Model Still

After seeing those results, I designed a follow-up, wrote its design into the lab file before it ran, and stored its results in results/ingest_followup.json. It turns early stopping off with early_stopping=False, so the trees learn from every row the job hands over and nothing is left to chance. The lab checked that the new trees give identical answers on a second fit and with another seed. Everything else is the same: faults, shares, seeds, test rows.

A line chart headed ingest_followup.json, designed after the results: early stopping off, the mean of the seeds, titled with the model held still, only the shuffle hurt every time. Accuracy from 0.76 to 0.86 against the share of files touched, 5, 10, 25 and 50 percent, with a dashed line for clean near 0.81. Sent twice, dashed, starts near 0.83 and settles near 0.81. Late files, dashed, stays near 0.81 to 0.82 and ends near 0.80. Column renamed, dashed, starts near 0.79, rises near 0.82, ends near 0.81. Rows shuffled, solid, starts near 0.80 and falls to about 0.78. Beneath: clean 0.8142; runs below clean: twice 10 of 20, late 12 of 20, shuffled 20 of 20, renamed 3 of 4; shuffled, every run, 0.7757 to 0.8014.

The clean trees now score 0.8142, and 0.7520 without the lag. With the model held still, the picture is clearer, and it is not the one I expected. The copies still went both ways, from 0.7855 to 0.8452, above clean in 10 of 20 runs. Late files went from 0.7705 to 0.8537, below clean in 12 of 20. The rename scored below clean in 3 of its 4 runs. Only the shuffle scored below clean in every one of its 20 runs, from 0.7757 to 0.8014. Its mean falls as the share of shuffled files grows.

How much did each fault change the model's answers? I measured this after the results, from the stored test counts. At the headline, the copies changed 10.5% of the test answers (955 of 9,063 rows). The late files changed 4.2%, the shuffle 3.7% and the rename 3.6%.

Copying 66 of 756 files changed one answer in ten. McNemar's test says the copied model beat the clean one by more than luck, with p far below 0.05. But across five seeds, the same fault made the model better in some runs and worse in others. The test sees that these two particular models differ. The seeds show that the direction depends on which files happened to be copied. The fault did not make models better. It made the model different, in a way nobody chose, and a team watching only accuracy would have called it an improvement.

The Fixes Give Back the Clean Table

Each fault has a fix that follows from the first half of the lesson. For copies and shuffles, keep one row per key, in key order, before building any input. For late files, rebuild the day once its file lands, which is a backfill. For the rename, map nsw_demand back to nswdemand; in a real pipeline, the schema check should stop the file until someone does.

Two panels headed ingest.json: the four fixes, every run, titled keyed, ordered and complete gives back the clean table. Keyed repair: 40 of 40 runs with copies or shuffles, equal to the clean table, value for value. Retrained on it: 0.8193, each of the four headline repairs; clean 0.8193. Beneath: the 20 backfills and 4 rename fixes equal the clean table by construction: they rebuild from every file or undo the rename.

The repairs are not all equally strong evidence. In the 40 runs with copies or shuffles, keeping one row per key in key order gave back a table equal to the clean one, value for value. The lab was written to stop if one was not. That is a real test, and it is paying off. Because every record has a key, you can repair copies and shuffles without knowing which files went wrong. Sorting by key puts every row back where it belongs.

The other 24 repairs equal the clean table by construction. The 20 backfills rebuild from every file, and the 4 rename fixes undo the rename. A late file or a renamed column cannot be repaired from the table alone: you need the missing file, or someone who knows what the new name means. Retrained, each of the four headline repairs scored 0.8193, the clean accuracy.

Try It Yourself

This script is the lab made small. It downloads the same data and sends it as one file a day. It applies each of the four faults to 10% of the files with seed 0, and trains the trees on each table. It prints what each table got and what the trees scored on the clean test rows. Then it applies the fix for copies and shuffles. It does not need a GPU.

A real screenshot of VS Code with ingest_demo.py open, showing the docstring that says what the script does and how to run it, the imports of numpy, pandas and scikit-learn, the lines that load Elec2, turn the label into 1 for UP, add the key and cut the training rows at 80%, the clean test rows with the lag of the half-hour before, the split into one file a day with 10% of files touched at seed 0, and the start of the job function, which reindexes the columns and fills missing values with 0; the rest of it, which takes the lag from the row above, is further down. Beneath: copy it from the box on the slide.

Before you run this lab. You need Python 3 and two libraries: pip install scikit-learn pandas. scikit-learn holds the model and the download, and brings NumPy with it; pandas holds the files and the table. The first run downloads Elec2 from OpenML (under 1 MB compressed), so it needs an internet connection once. scikit-learn keeps a copy in a folder in your home directory (scikit_learn_data) for later runs.

It trains seven small models, so it is quick; I did not time it on a quiet machine. I ran it with scikit-learn 1.9.1 on a Mac. scikit-learn runs the same way on Windows and Linux, but I have not checked the numbers there. Another version of scikit-learn may give different decimals, so the first line printed is the version.

"""Four quiet ingestion faults, and what each does to training.

Data ingestion pipelines, made small. Elec2 arrives as one file
a day. A fault touches 10% of the files: sent twice, late,
shuffled, or a column renamed. For each case it prints the rows
the training table got, the copies, the missing rows, the wrong
lag inputs, and the trees' accuracy on clean test rows. Then the
fix for copies and shuffles: one row per key, in key order.

It needs Python 3 with scikit-learn and pandas
(pip install scikit-learn pandas). The first run downloads
Elec2 from OpenML (under 1 MB) and keeps a copy.
    python ingest_demo.py

Author: Roni Das
Created: 2026-09-30
"""
import numpy as np
import pandas as pd
import sklearn
from sklearn.datasets import fetch_openml
from sklearn.ensemble import HistGradientBoostingClassifier

# 45,312 half-hours, 48 a day. The row number is the key.
data = fetch_openml("electricity", version=1, as_frame=True,
                    parser="auto").frame
COLS = list(data.columns[:8])
src = data[COLS].astype(float)
src["label"] = (data["class"].astype(str) == "UP").astype(int)
src["key"] = np.arange(len(src))
cut = int(0.8 * len(src))  # learn from the first 80%
y = src["label"].to_numpy()
prev = np.r_[y[0], y[:-1]]  # the label of the half-hour before
test = src[COLS].assign(lag=prev)[cut:]  # clean test rows
print(f"scikit-learn {sklearn.__version__}")

# The source sends one file a day. A fault touches 10% of them.
days = src["key"][:cut] // 48
files = [f for _, f in src[:cut].groupby(days)]
hit = np.random.default_rng(0).random(len(files)) < 0.10


def job(landed):
    # The feature job. It trusts the order the rows landed in.
    t = pd.concat(landed).reindex(columns=COLS + ["label", "key"])
    t = t.fillna(0)  # a missing value becomes 0
    lag = t["label"].shift(1).fillna(t["label"].iloc[0])
    return t, t[COLS].assign(lag=lag)


def train(name, landed):
    t, X = job(landed)
    trees = HistGradientBoostingClassifier(random_state=0)
    trees.fit(X, t["label"])
    acc = (trees.predict(test) == y[cut:]).mean()
    copies = t["key"].duplicated().sum()
    gone = cut - t["key"].nunique()
    wrong = (X["lag"].to_numpy() != prev[t["key"]]).sum()
    print(f"{name:<16}{len(t):>7,}{copies:>7,}{gone:>7,}"
          f"{wrong:>6,}{acc:>8.4f}")


k = round(0.10 * len(files))  # the files after the rename
renamed = [f.rename(columns={"nswdemand": "nsw_demand"})
           for f in files[-k:]]
cases = {
    "clean": files,
    "twice": [g for f, h in zip(files, hit) for g in [f] * (1 + h)],
    "late": [f for f, h in zip(files, hit) if not h],
    "shuffled": [f.sample(frac=1, random_state=d) if h else f
                 for d, (f, h) in enumerate(zip(files, hit))],
    "renamed": files[:-k] + renamed,
}
print(f"{'case':<16}{'rows':>7}{'copies':>7}{'gone':>7}"
      f"{'lag':>6}{'acc':>8}")
for name, landed in cases.items():
    train(name, landed)

# No error, no copy, no gap: the renamed column just went empty.
raw = pd.concat(cases["renamed"]).reindex(columns=COLS)
print(f"rows with no nswdemand: {raw['nswdemand'].isna().sum():,}")

# The fix for copies and shuffles: one row per key, in order.
for name in ("twice", "shuffled"):
    rows = pd.concat(cases[name]).drop_duplicates("key")
    train(name + ", fixed", [rows.sort_values("key")])

The Lab Report

A real terminal recording of python ingest_report.py. It opens with Elec2 in 756 daily files, 36,249 training rows and 9,063 clean test rows, and 65 checks against a rebuild written apart from the lab: all agree. Then five sections. 1, the table each fault built at 10% of files and seed 0: files, rows, copies, missing, lag off and no demand for clean and the four faults. 2, the checks on the landed files: count, copies, missing, order, schema and nulls, yes or a dash for each case. 3, the model with default trees: clean 0.8193, random_state 1 to 9 from 0.7884 to 0.8322, without the lag 0.7527, then each fault's accuracy, McNemar p, range over all runs and runs below clean, and a note that the rename keeps the clean split and that every model used all 100 rounds. 4, the follow-up with early stopping off, designed after the results: clean 0.8142, then accuracy, McNemar p, share of answers changed and runs below clean. 5, the fixes: in the 40 runs with copies or shuffles the keyed repair equals the clean table, a real test; the 20 backfills and 4 rename fixes equal it by construction; each headline repair scores 0.8193. Beneath: the report rebuilds the tables with its own code, refits the headline models, and stops on the first number that does not come back.

The report lives in scripts/labs/dataeng/ingest_report.py. It reads the lab's two stored files and the Elec2 data from scikit-learn's local copy. It does not trust the lab's code. With its own, separately written code, it lands the files, builds the tables and runs the six checks again. It refits the headline models, and stops on the first number that does not come back exactly. It makes 65 checks, and they all agree. It changes nothing in the lab's files.

Its json mode writes the numbers the figures read to results/ip-report.json. The demo mode checks the student script's stored run. The box mode writes the playground on the next slide and checks it against the lab: 182 checks, at every share and seed the box can rebuild.

What came before the run, in ingest_lab.py: the data, the source and the job, and the model with its two references. Also the four faults, the shares and seeds, the six checks, the fixes, and McNemar's test for the headline runs. What came after I saw the main results: the follow-up with early stopping off, designed and written down before it ran. What came after all the results, in the report: the share of answers that changed, the row counts per file, and the sketched first files.

Run the Checks Yourself

This box has no model in it. It holds the real labels of all 36,249 training half-hours, and which files each fault touched at every share and seed the lab ran. It also holds the landing order of every file shuffled at 10% with seed 0. From those it rebuilds each landed table and runs the checks, in your browser. The box holds keys, not columns, so its two column checks both read a mark on each renamed file, not real column names. The report checked that every row count, copy count, missing count, share of wrong lag inputs, share of UP rows and check result it gives matches results/ingest.json.

As it is, the box prints the headline table. For each case it shows the rows, copies and missing keys, the share of wrong lag inputs, the share of rows that are UP, and which checks fired. The clean files fire none, and each fault fires its own.

Then try other doses: table(0.50, 3) shows every fault except the shuffle at half the files with seed 3 (the box holds the shuffled order for the headline only). To see copies arrive, print the first key of each landed file with [f[0] for f, _ in land('twice')[:6]]. Try fix('twice') == land('clean') and fix('late') == land('clean'). Last, a question to answer before you run it: at table(0.05, 0), which checks will fire for the late files, and which will not?

The Code, Part by Part

Loading. fetch_openml("electricity", version=1) downloads Elec2 once and reads the local copy after that. The first eight columns are the inputs; label is 1 for UP; key is the row number, which names each half-hour. cut is 36,249: the trees learn from the rows before it.

The clean test rows. prev shifts the labels down by one row, so each row holds the label of the half-hour before. test is the last 9,063 rows with that lag added. They are built from the clean source, never from the landed files.

The files. files splits the training rows into one table per day, by key // 48. hit marks the files a fault touches: 10% of them, chosen with seed 0.

job. This is the feature job, and it is where the faults do their harm. pd.concat stacks the landed files in landing order. reindex picks the expected columns by name, and silently gives an empty value for a column that is not there; fillna(0) then turns it into 0. The lag is shift(1): the label of the row above, which is the half-hour before only when the rows are in time order.

How to Build Ingestion That Cannot Go Wrong Quietly

A hand-sketched column of six boxes joined by arrows, headed ingestion that cannot go wrong quietly, titled land, key, count, check, order, backfill. 1, land every file as it came; record when it happened and when it arrived. 2, give every record a key; write with a keyed upsert. 3, compare each file's row count with the source's count. 4, check each file's column names before the table. 5, build inputs like the lag in key order, never landing order. 6, rebuild a day when its late file lands: a backfill. Beneath: here, checks written for these four faults fired in 64 of 64 runs; accuracy moved both ways.

Land every file as it came. Keep the raw copy, and record two times with every record: when it happened and when it arrived. Build time windows from the first. The raw copy is what every later fix is built from.

Give every record a key, and write with keyed upserts. Then a retry, a replay or a backfill cannot create a copy, and a version column decides which copy wins.

Count at the door. Ask the source for a count per file or per day; many exports can write one next to the data. Compare it with what landed. In the lab this caught the copies and the gaps in every run.

Check the columns before the table. Compare every file's column names and types with what you expect, and stop the file, or set it aside, on a mismatch. Never let a missing column quietly become 0.

Build order-dependent inputs in key order. Any input that looks at "the row before", a running total or a time window must sort by event time or key first. Never trust landing order.

Backfill late days. When a late file lands, rebuild the days it belongs to, and keep recent days marked as provisional until the wait is over.

A Nightly Pull, or Change Data Capture?

A two-column page headed choosing how to ingest, titled a nightly pull, or change data capture? Left, a nightly pull fits when: the model decides once a day or less; the source is an outside API you cannot install on; deleted rows do not matter, or the source marks them; a day-old value costs little. Right, CDC or streaming fits when: a decision needs changes from minutes ago; deletes must reach the model; the database must not take heavy reads; someone can run and watch a broker. Beneath, left: either way, keys, counts and column checks. Beneath, right: a retry here made 39,417 rows of 36,249.

Use a nightly batch pull when the model decides once a day or less. A day-old value costs little there, and a pull is the simplest thing to run and to debug.

Use a pull when the source is an outside API you cannot install anything on. Poll it on a schedule, a page at a time, within its .

Use a pull only when deletes do not matter, or the source marks them. A watermark on updated_at cannot see a row that is gone, so many teams ask the source to mark a row as deleted instead of removing it.

Do not use a pull when a decision needs changes from minutes ago. Use CDC for a database and streaming for app events.

Do not use a pull when deletes must reach the model, such as a cancelled order or an account whose owner asked to be forgotten. CDC sees deletes; a watermark does not.

Do not run CDC without someone who can run and watch a broker and its slots. An abandoned slot can fill a database's disk. And whichever mode you choose, the lab's lesson holds: keys, counts and column checks. Every mode can retry, fall behind or meet a renamed column.

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 dataset: Elec2, 1996 to 1998; one model: trees with a lag input; faults I built on purpose; the design came first; early stopping off: after; five seeds a share of files. They are not: not a rate for other data; not every kind of model; not a bug found in the wild; not changed after the run; designed knowing the results; not a law about direction.

One dataset, one model. Everything here is one electricity market from 1996 to 1998 and one kind of model, boosted trees with a lag input. With other data and other models, what each fault costs would move. The lag input makes this model unusually sensitive to order, which makes the shuffle a clear demonstration and a poor estimate for models without such an input.

Faults I built on purpose. The copies, gaps, shuffles and rename are simulations of failures that happen in real pipelines. They are not a bug I found in someone's system. A real fault can be messier: copies of only some rows, or a rename that reaches only some sources.

Five seeds, not a law. Five seeds per share show the spread; they do not prove a direction. Where I say a fault went both ways, that is five runs, here.

The follow-up came after. Turning early stopping off was designed after I saw the main results, and written down before it ran. The main run's McNemar tests were declared in advance. As the model slide explains, they cannot separate the fault from the change in which rows early stopping set aside.

The checks were built for these faults. I wrote each check together with the fault it matches, so their firing in 64 of 64 runs was expected. "None fired on the clean table" rests on one clean table. The missing-keys check works only because the keys here are row numbers with no gaps. A random id, or an order_id with holes of its own, cannot show a missing record. And a fault none of them was built for would pass all six. If the source began sending demand in kilowatts instead of megawatts, every file would still have the right rows, keys, order and columns.

So do not read this lab as "six cheap checks catch every ingestion fault". They catch these four, and a pipeline needs others too, such as checks on the range and spread of each column. The checks are also strict: any copy, any gap, any key out of order. A real pipeline receives some data out of order for good reasons, so its order check would compare event times within a window. And its row count needs the source to send a count at all.

What to Do Next

A hand-drawn list headed before you trust a training table, titled five questions for every ingestion job. Key?: does every record have a key, and is every write keyed on it? Count?: does anything compare the rows landed with the rows sent? Order?: is any input built from the order the rows happened to land in? Columns?: does a missing column stop the file, or quietly become 0? Late?: when a late file lands, does anything rebuild its day? Beneath: here, a job that trusted landing order got 1,406 lag inputs wrong from 66 shuffled files.

Take one ingestion job your team runs and ask the five questions on the card. If the answer to "key?" is no, start there: a key per record and keyed writes make every later fix possible. If "count?" is no, ask the source team whether their export can write a row count next to each file. It is a small change, and it makes the most common faults visible.

Then run a small version of this lab on your own data. Take one day's landed files, and count the rows per file against what the source says it sent. Then look for repeated keys, and check the key order and the column names. If all four are clean, you have learned something useful about your pipeline. If one is not, you found it before a model did.

The next lesson, ETL vs ELT, is about where the cleaning happens once the data has landed: before it is loaded into the warehouse, or after.

A closing card headed to keep, titled ingestion faults are quiet; check the table, not the score. In large type: 36,249 records, 39,417 rows. Beneath: 66 files sent twice; no error, and in this run accuracy even rose, 0.8193 to 0.8230. Then: rows shuffled in 5% to 50% of files: below clean in 20 of 20 runs, and in 20 of 20 with the model held still. Then: checks written for these four faults caught each of them in every run; a change of units would pass them all. Last: one dataset, faults I built: a way to check, not a law.

If you keep one thing from this lesson, keep the gap between the two numbers at the top of the card. Those 3,168 copies raised no error, and in that run they even raised the score. A pipeline can be wrong in a way its model rewards, so judge the table by its keys, counts, order and columns, not by the accuracy it produces.

Knowledge Check

Knowledge Check

4 questions - Score 80% to pass

Q1

In the lab, 66 of 756 daily files were delivered twice and the job appended both copies. What did the accuracy of trees trained on that table do on the clean test rows?

Q2

Shuffling the rows inside 66 files left the training table with exactly 36,249 rows. Which of the lab's checks caught it?

Q3

With default settings, the same trees trained on the same clean table with random_state 1 to 9 scored 0.7884 to 0.8322. What does that mean for watching accuracy to catch ingestion faults?

Q4

Why does writing with keyed upserts make it safe to replay a whole topic or rebuild from the landing zone?

A flowchart headed the path in this lesson's lab, titled where each fault gets in. The source: one file a day, leads to the connector lands each file, which leads to a cylinder, landing zone: files as they arrived, which leads to the feature job builds the training table, which leads to the trees learn from the table. A separate box, clean test rows, from the live path, joins the trees with a dotted line labelled score it. Beneath: copies, gaps and shuffles get in at the connector; the rename gets in at the source; nothing touches the test rows.

The rows the model is tested on never pass through the faulty path. They come from the clean source, the way a correct live system would give them. Only the table the model learns from is broken. That keeps the question simple: when the training table is wrong in one quiet way, what changes?

offset

Most platforms settle on a simple rule. Pull from outside APIs and slow tables, where a schedule is fine and you cannot install anything. Push from the busy databases and event streams you need close to real time. The rest of the first half makes that choice concrete.

Best for
Batch pullUp to the scheduleA read each runNo, with a watermarkOutside APIs, slow tables, daily reports
CDCClose to real timeSmall: it reads the logYesFeatures kept in step with a busy database
Event streamingClose to real timeA broker to runOnly if the app sends oneClicks and actions, live features

Mature platforms often run all three at once, one per kind of source. The beginner's mistake is to force everything through a nightly batch, then wonder why a model cannot react to something that unfolds in ten minutes. Freshness that ingestion never captured cannot be added later.

before
before

The event lands in . A consumer pulls it, writes it into the landing zone, and then commits its offset. Committing the offset is the backbone of reliable ingestion. If the consumer crashes after writing but before committing, it restarts from the last committed offset and reads those events again. Nothing is lost, but something may be delivered twice. Kafka's documentation puts it plainly: Kafka "guarantees at-least-once delivery by default", and Debezium's documentation says it "provides at least once delivery" after a fault. So the write into the landing zone must be idempotent, and the duplicates slide and the lab are both about what happens when it is not.

table.include.list limits it to the two tables you care about. Reading every table of a busy database is a good way to fill a Kafka cluster you did not plan for. topic.prefix has no default and must be set; it names the Kafka topics.

decimal.handling.mode is the kind of setting that quietly decides whether money is right. The default, precise, keeps amounts exact but sends them as encoded bytes that many readers cannot decode. string keeps the digits as readable text. double turns them into floating-point numbers, which Debezium's documentation warns "might result in a loss of precision", so avoid it for money.

One more setting protects the database itself. A replication slot that no reader drains makes Postgres keep its log files, because it cannot throw away log it may still owe that slot. By default there is no limit: with max_slot_wal_keep_size at its default of -1, Postgres's documentation says slots "may retain an unlimited amount of WAL files". The disk fills, and a database whose log disk is full can shut down. From Postgres 13 you can set that cap, and a slot that falls further behind than the cap is marked invalid instead of taking the database down. Either way, if you switch a CDC connector off, drop its slot.

Choose the version with care. A timestamp like updated_at is tempting, but two different updates can carry the same timestamp, and the clocks of different machines drift apart. Under a strict "newer wins" rule, a tied update is silently dropped. A better version only ever goes up and never repeats. In Postgres that is the log sequence number, or LSN: its position in the write-ahead log, which the documentation says is "increasing monotonically with each new record". Debezium puts it in every event, as source.lsn. An application can also keep its own row version counter.

Deletes need the same care. If a delete simply removes the key, a replayed old insert finds no row and writes it back, and the deleted order comes back to life. So keep a delete marker under the key, often called a tombstone, with the version of the delete. It then wins against older copies like any other write.

Here is the whole rule as runnable code, on a tiny made-up stream with a version that only goes up and one delete. Watch the last two lines. Replaying the stream changes nothing, and without delete markers the deleted order comes back.

The payoff is large. With keyed upserts, you can replay a whole Kafka topic or rebuild everything from the landing zone, and land in exactly the same final state. That one property, safe replay, makes the rest of ingestion recoverable. It is also why the backfills on the next slide are safe to run while the normal pipeline keeps flowing. In the lab it is also the fix: keeping one row per key gave back the clean table exactly.

New column that may be empty
Safe
Add it, if schema merging is turned on (Spark's mergeSchema, Delta Lake's schema evolution). By default Delta Lake refuses the new column
Column made able to be emptySafe to storeWatch its share of empty values before any fill. A fill of 0 hides them, as it hid 3,609 in the lab
Column renamedBreakingSet the rows aside and alert a person. A rename reads as one column missing and a new one appearing
Column type changed, number to textBreakingSet aside, alert, and do not convert blindly
Column droppedBreakingAlert, and decide on purpose what to do downstream

The mechanism behind "set aside and alert" is the dead-letter queue: a side channel for records that fail to parse or fail a check. When a record fails, it must not quietly vanish, and it must not crash the whole run either.

A flowchart headed the file that fails a check, titled set it aside with its reason; keep the rest flowing. A box, a file arrives, leads to a diamond, columns as expected? Yes leads to a cylinder, the trusted table. No leads to a cylinder, set aside: the file and the reason, which leads to a box, a person is told, and from there a dotted line labelled the mapping is fixed goes back to the diamond. Beneath: in the lab the schema check fired on the first renamed file, the one for 18 March 1998; a side channel like this is called a dead-letter queue.

The failing record goes to the side channel with the reason it failed. There it can be fixed and replayed once the mapping is corrected, and a person is told when the rate of failures rises. Connect has a dead-letter queue built in, with a catch. It is for sink connectors, the ones that write data out of Kafka. And it is useful only when errors.tolerance is set to all: under the default, none, the task stops at the first bad record.

The principle is simple: absorb safe changes automatically, and never write breaking ones blindly. A crashed pipeline is annoying but honest. A pipeline that quietly corrupts a feature for a week is the expensive kind of failure. You find out only when the model has already made a week of predictions on bad data.

A two-column page headed the four faults, fixed before the run, titled what each one does, and where it comes from. Left, in the lab; right, in a real pipeline. Sent twice: a touched file lands twice, the copy right after it; a retry after a lost acknowledgement, and a job that appends. Late: a touched file is missing when the table is built; an export that arrives after the nightly build. Shuffled: a touched file's rows land in a random order; the worst case of rows read back from several partitions. Renamed: the last 10% of files call nswdemand nsw_demand; a change upstream that nobody announced. Beneath, left: only the training table is broken. Beneath, right: the test rows are always clean.

The four faults each touch a share of the files: 5%, 10%, 25% or 50%. For three of them, the touched files are chosen at random, with five different seeds: the starting numbers for the random choice. That shows how much depends on which files were hit. The rename has no randomness: it always hits the most recent files, because a rename changes everything after it. The headline cases on the next slides are 10% of files with seed 0. The test rows are never faulty.

A row count sees only two of the four faults. The shuffled and renamed tables have exactly 36,249 rows, the right number, so the count check passed them both. Only a check on key order sees the shuffle, and only a check on column names or empty columns sees the rename.

The declared test was McNemar's test. It counts only the test rows where exactly one of two models is right, and asks how likely a split that uneven would be by chance alone. Its answer is a p value: the chance of a split at least that uneven if the two models were equally good. By the usual rule, a p below 0.05 means luck is an unlikely explanation. Here p went from 0.273 for the copies to far below 0.05 for the late files.

It answers a narrower question than it seems to. It compares one model with another, and for the copies, the late files and the shuffle, the two models also differ in which rows early stopping set aside. So for those three it cannot say how much of a gap the fault caused, and I do not read it as saying so. The rename is the exception. It keeps the clean rows with their labels in the clean order, so its trees set aside exactly the same rows. Its p of 0.0136 is the one comparison here that the split does not muddle.

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

A real screenshot of VS Code's terminal after running python ingest_demo.py. It prints scikit-learn 1.9.1, then a table with the columns case, rows, copies, gone, lag and acc. Clean: 36,249, 0, 0, 0, 0.8193. Twice: 39,417, 3,168, 0, 27, 0.8230. Late: 33,081, 0, 3,168, 27, 0.7849. Shuffled: 36,249, 0, 0, 1,406, 0.8143. Renamed: 36,249, 0, 0, 0, 0.8134. Then: rows with no nswdemand: 3,609. Then twice, fixed: 36,249, 0, 0, 0, 0.8193; shuffled, fixed: 36,249, 0, 0, 0, 0.8193.

When I ran it, it printed scikit-learn 1.9.1 and then the table. Clean: 36,249 rows and 0.8193. Twice: 39,417 rows, 3,168 copies, 27 wrong lag inputs and 0.8230. Late: 33,081 rows, 3,168 gone, 27 wrong lags and 0.7849. Shuffled: 36,249 rows, 1,406 wrong lags and 0.8143. Renamed: 36,249 rows and 0.8134, with 3,609 rows that had no nswdemand. Both fixed tables printed 36,249 rows and 0.8193. All of it matches results/ingest.json, the longest printed line was 51 characters, and the report's demo mode checks every number.

To see a heavier dose, change 0.10 in the hit line to 0.50. The shuffled case then gets far more wrong lag inputs. The lab's stored run at that share with seed 0 had 20.4% of its lag inputs wrong, and scored 0.7811.

train. Trains the trees with seed 0 on the table and scores them on the clean test rows. It prints the rows, the copies (duplicated), the missing keys, the wrong lag inputs and the accuracy. A lag input is wrong when it differs from prev for that row's key.

The cases. Each case is a list of files in landing order. twice repeats a touched file. late leaves it out. shuffled puts its rows in a random order with sample(frac=1). renamed gives the last 10% of the files a new name for one column.

The fix. drop_duplicates("key") keeps one row per key and sort_values("key") puts the rows back in time order. Fed to the same job, both broken tables become the clean one.

In the lab file, ingest_lab.py does the same for every share and seed, and runs the six checks. It repairs every table, checks that it equals the clean one, and stores it all in results/ingest.json. Its followup mode repeats the grid with early stopping off.