Data Engineering for AI

ETL for machine learning

ETL moves data from where it lands to where a model can use it, and the job that matters is making that move safe to run twice.

Read these first

On this page 8
  1. The three letters
  2. The property that matters most
  3. How to get it
  4. The gate before the shelf
  5. Somewhere you have seen this
  6. What is honestly hard here
  7. Remember this
  8. What to learn next

One lesson, three depths. Pick the one that fits you today — you can switch any time.

Beginner — No maths. Plain English.

ETL means extract, transform, load — take data from where it lands, reshape it, and put it where a model can use it.

Think about a kirana shop receiving a delivery. Boxes arrive at the back door. Someone opens them, checks nothing is broken, and puts each item on the right shelf.

Nobody sells from the boxes. The unpacking step is what makes the shop usable.

ETL is that unpacking step, for data.

The three letters

Extract. Pull data out of wherever it lives — a database, a payments provider, a folder of files somebody uploads.

Transform. Fix it up. Correct types, standardise spellings, drop repeats, calculate the columns a model needs.

Load. Write the result somewhere the training code and the live system can both read.

You will also meet ELT, where you load the raw data first and transform it afterwards inside the warehouse. Same three steps, different order, and it suits modern warehouses better because storage got cheap.

The property that matters most

Pipelines fail. The network drops, a server restarts, a file arrives late. So somebody re-runs yesterday's job.

If re-running adds everything a second time, you now have every row twice. The model then trains on a table that lies. Fixing that by hand at 2am is somebody's actual job tonight.

A job that gives the same result whether you run it once or five times is called idempotent. Making your loads idempotent is the single highest-value habit in this whole section.

How to get it

Do not append blindly. Give every record a key that identifies the real thing — a payment id, an order number — and write with that key.

If the key is already present, replace the row. If not, add it. Running the job again then changes nothing.

   APPEND                            UPSERT (update or insert)
   run once   ->  5 rows             run once   ->  5 rows
   run again  -> 10 rows             run again  ->  5 rows
   run again  -> 15 rows             run again  ->  5 rows

The gate before the shelf

Before anything is written, check it. Are the columns the ones you expected? Are the numbers numbers? Are there negative prices?

If a check fails, stop. Do not load a broken batch and hope somebody notices.

A pipeline that fails loudly is doing its job. A pipeline that quietly loads bad data is worse than no pipeline. That bad data now carries your organisation's trust.

Somewhere you have seen this

Your bank statement arriving each month, in the same format, with the same columns. Behind it, a job pulls transactions, tidies them and writes them out, every day, without a person involved.

What is honestly hard here

Upstream systems change without telling you. A field that was always filled starts arriving blank. A vendor renames a column in a routine update.

You cannot prevent this. You can only detect it fast, which is why the checks matter more than the transformations.

Backfilling is hard too. Fix a bug today and you must decide whether to recompute two years of history. Those old numbers may have driven decisions people still rely on.

Remember this

  • ETL is extract, transform, load. ELT does the same work in a different order.
  • Idempotent means running twice is the same as running once. Use keys, not appends.
  • Check before you load. A pipeline that stops loudly beats one that loads quietly.

What to learn next

Developer — Code and libraries.

Here is the difference between an append and an upsert, measured on a replay — which is the situation that actually breaks pipelines.

Setup

bash
pip install pandas

A pipeline that survives being run twice

etl.py
import pandas as pd
from pandas.api import types as pt

SCHEMA = {"txn_id": pt.is_integer_dtype, "customer": pt.is_string_dtype,
          "amount": pt.is_float_dtype,   "ts": pt.is_datetime64_any_dtype}

def extract(day):
    batches = {
        "2026-08-18": [(101, "c1", 250.0), (102, "c2", 990.0)],
        "2026-08-19": [(103, "c1", 120.0), (104, "c3", 75.5)],
        "2026-08-20": [(102, "c2", 990.0), (105, "c2", 40.0)],   # 102 is re-sent
    }
    return pd.DataFrame(batches[day], columns=["txn_id","customer","amount"]).assign(ts=pd.Timestamp(day))

def validate(df):
    problems = []
    for col, ok in SCHEMA.items():
        if col not in df.columns: problems.append(f"missing column {col}")
        elif not ok(df[col]):     problems.append(f"{col} has wrong type: {df[col].dtype}")
    if "amount" in df and pt.is_numeric_dtype(df["amount"]) and df["amount"].lt(0).any():
        problems.append("negative amount")
    if "txn_id" in df and df["txn_id"].duplicated().any():
        problems.append("duplicate txn_id inside one batch")
    if problems: raise ValueError("batch rejected: " + "; ".join(problems))
    return df

def load_append(wh, batch): return pd.concat([wh, batch], ignore_index=True)
def load_upsert(wh, batch, key="txn_id"):
    return pd.concat([wh, batch], ignore_index=True).drop_duplicates(subset=key, keep="last").reset_index(drop=True)

days = ["2026-08-18","2026-08-19","2026-08-20"]
empty = extract(days[0]).iloc[0:0]          # right columns, right types, no rows
append_wh, upsert_wh = empty.copy(), empty.copy()
for run in ("first run", "operations replayed the week"):
    for d in days:
        b = validate(extract(d))
        append_wh = load_append(append_wh, b)
        upsert_wh = load_upsert(upsert_wh, b)
    print(f"{run:32s} append: {len(append_wh):2d} rows   upsert: {len(upsert_wh):2d} rows")
print()
print("copies of txn 102 in the append warehouse:", int((append_wh['txn_id']==102).sum()))
print("copies of txn 102 in the upsert warehouse:", int((upsert_wh['txn_id']==102).sum()))
print()
print(upsert_wh.sort_values("txn_id").to_string(index=False))
print()
bad = extract("2026-08-20").assign(amount=lambda d: d["amount"].astype(str))
try: validate(bad)
except ValueError as e: print("gate fired ->", e)
Output
first run                        append:  6 rows   upsert:  5 rows
operations replayed the week     append: 12 rows   upsert:  5 rows

copies of txn 102 in the append warehouse: 4
copies of txn 102 in the upsert warehouse: 1

 txn_id customer  amount         ts
    101       c1   250.0 2026-08-18
    102       c2   990.0 2026-08-20
    103       c1   120.0 2026-08-19
    104       c3    75.5 2026-08-19
    105       c2    40.0 2026-08-20

gate fired -> batch rejected: amount has wrong type: object

What the counts say

The append warehouse doubled on a replay. 6 rows became 12. Every downstream average, count and rate is now wrong by a factor that depends on how many times operations retried.

The upsert warehouse did not move. 5 rows, then 5 rows. That is idempotency, and it took one drop_duplicates(subset="txn_id", keep="last").

Notice txn 102 appeared four times in the append store. It was re-sent by the source on the 20th, then the whole week was replayed. Two independent duplication sources compounded. This is normal, and it is why "the source will not send duplicates" is not a design.

The gate rejected a batch with the wrong type before a single row was written. amount arrived as text, which would silently break every arithmetic operation downstream. Catching it at the door costs one function.

Why keep="last" and not keep="first"

A re-sent record is often a correction. The source found an error and sent the row again with a fixed amount.

keep="last" treats the newest arrival as authoritative. If your source instead sends the original first and duplicates by accident, keep="first" is right. Decide deliberately, and write down which you chose — this is a business rule wearing the costume of a parameter.

For real correctness you want a sequence number or an updated timestamp from the source, sort by it, and keep the last. Arrival order is not reliable.

Making a real database load idempotent

The pandas version above demonstrates the idea. In a warehouse you would write:

sql
-- PostgreSQL / SQLite
INSERT INTO transactions (txn_id, customer, amount, ts)
VALUES (?, ?, ?, ?)
ON CONFLICT (txn_id) DO UPDATE
SET customer = excluded.customer,
    amount   = excluded.amount,
    ts       = excluded.ts;

Two other patterns worth knowing:

Partition overwrite. Write each day into its own partition and replace that whole partition on re-run. Idempotent by construction, and the standard approach on object storage.

MERGE. The SQL standard's upsert, supported by most warehouses and by Delta Lake and Iceberg. More expressive than ON CONFLICT, and it lets you handle deletes in the same statement.

Common mistakes

Using INSERT and relying on operations never to re-run a job. They will re-run it. That is what operations does when something fails.

Making the key (everything). If your key is all columns, a corrected row is a new row and you keep both. The key must identify the real-world entity, not the bytes.

Validating after loading. The check must run before the write. Otherwise the gate is a report about damage already done.

Silent coercion. pd.to_numeric(x, errors="coerce") turns anything unparseable into NaN without complaint. Convenient, and it hides exactly the upstream change you needed to hear about. Use errors="raise" in a pipeline, and handle the exception deliberately.

No row-count logging. Record rows in and rows out at every stage. A step that usually drops 2% and today dropped 60% is an incident, and this is the only way you will see it.

Transforming before validating. Your transformation code assumes the schema. Check the schema first, or your error messages will point at the wrong line.

Try it yourself

Change keep="last" to keep="first" and re-run. The row for txn 102 keeps its original timestamp of the 18th rather than the 20th. Both behaviours are defensible; only one matches your source's semantics.

Then add a check to validate requiring ts to be within the last 7 days, and feed it a batch dated 2020. Confirm the gate fires. Staleness checks catch a whole family of upstream failures that type checks miss.

What to learn next

Researcher — Mathematics and papers.

Delivery semantics and where idempotency comes from

Distributed pipelines offer three delivery guarantees:

  • At-most-once. No retries. Messages can be lost.
  • At-least-once. Retry until acknowledged. Duplicates are guaranteed, not only possible.
  • Exactly-once. No loss, no duplicates.

Exactly-once message delivery is impossible in an asynchronous network with failures — this follows from the Two Generals problem, and relatedly from the FLP impossibility result (Fischer, Lynch and Paterson, 1985) on deterministic consensus with one faulty process.

What systems actually provide is exactly-once processing semantics: at-least-once delivery combined with idempotent or transactional application of effects. Kafka's exactly-once support is built from idempotent producers (sequence numbers per producer and partition, deduplicated by the broker) plus transactional writes spanning consume, process and produce.

The practical implication: idempotency is not a nicety you add for tidiness. It is the mechanism by which the guarantee you want is constructed from the guarantee you can have.

Formal statement

A load operation $L$ over a warehouse state $W$ and a batch $B$ is idempotent if:

$$ L\big(L(W, B),\, B\big) \;=\; L(W, B) $$

Upsert keyed on a primary key satisfies this. Append does not. Partition overwrite satisfies a stronger property — the result depends only on the latest input for that partition, independent of history, which makes backfills safe by construction.

A stronger and more useful property is commutativity plus idempotency, which makes the operation a join-semilattice element and therefore order-independent. Last-write-wins upsert keyed by a monotonic source sequence number has this; upsert keyed by arrival order does not, which is why out-of-order replay can corrupt a warehouse that looks correct under sequential replay.

ETL against ELT

The shift from ETL to ELT is a consequence of the storage/compute cost curve, not of fashion.

ETL transforms before loading. It suits expensive storage and rigid schemas, and it means the warehouse holds only modelled data. The cost is that the raw form is discarded, so a transformation bug is unrecoverable without re-extraction from the source.

ELT loads raw, then transforms with the warehouse's own compute. Advantages: raw data is retained so any transformation can be recomputed; transformations are SQL and therefore versionable, testable and reviewable; and the warehouse's optimiser does the work. This is the dbt model, and its real contribution was bringing software engineering practice — version control, tests, dependency graphs, documentation — to transformation logic that had lived in untracked stored procedures.

The medallion pattern (bronze = raw, silver = cleaned and conformed, gold = aggregated for consumption) is the same idea with names attached. Its value is that each layer's contract is explicit and each is independently recomputable from the one before.

Orchestration semantics

Airflow, Dagster and Prefect all model pipelines as DAGs, and differ in what a node represents.

Airflow's unit is a task: a thing that runs. Dagster's is a software-defined asset: a thing that exists, with the computation that produces it attached. The asset framing makes lineage and staleness first-class — you can ask which assets are out of date with respect to their inputs — and matches ML work better, where the deliverable is a table or a model rather than a script execution.

Whatever the tool, three properties are non-negotiable:

  1. Deterministic partitioning. Every run is parameterised by a partition key, usually a time interval. SELECT ... WHERE ts >= '{{ ds }}' and never WHERE ts >= NOW() - INTERVAL '1 day', because the second is unreproducible.
  2. Retry safety. Every task must be safe to retry, which reduces to the idempotency condition above.
  3. Backfill as a first-class operation. Recomputing 2019 must use 2019's parameters, and either 2019's code or an explicitly recorded decision to use today's.

Data contracts

The failure mode ETL cannot engineer away is upstream schema change. Data contracts move the schema from an implicit assumption to an enforced interface: the producer commits to a schema, and the consumer's expectations are checked at the boundary.

Mechanisms in production use:

  • Schema registries (Avro, Protobuf) with declared compatibility modes. Backward-compatible evolution — adding an optional field with a default — is permitted; removing a field or narrowing a type is rejected at publish time rather than discovered at 3am.
  • Assertion-based validation — Deequ, Great Expectations, TFDV — running the checks as a gate rather than as a report.
  • Circuit breakers. A failed contract halts the pipeline. The alternative, loading anyway and alerting, means bad data reaches consumers before the alert is read.

Breck et al. (2017), The ML Test Score, IEEE Big Data, scores production ML systems across four categories. The data tests — feature distributions, schema conformance, privacy checks — are where most systems score lowest, and they are the cheapest points to earn.

Cost and scale notes

  • Columnar storage with predicate pushdown reduces scanned bytes by roughly the selectivity of the filter multiplied by the fraction of columns read. Partition pruning is the highest-leverage optimisation available; it is a design decision at write time, not a query hint.
  • Small-file proliferation is the standard failure of streaming ingestion into object stores. Metadata operations dominate, and compaction becomes mandatory rather than optional.
  • Incremental models — recomputing only changed partitions — cut cost linearly in the changed fraction, at the price of needing a correct change-detection predicate. A wrong predicate produces silently stale data, which is worse than a slow full refresh.

Reading

  • Kleppmann, Designing Data-Intensive Applications, O'Reilly 2017 — chapters 10 and 11 are the definitive treatment of batch and stream processing.
  • Fischer, Lynch and Paterson, Impossibility of Distributed Consensus with One Faulty Process, JACM 1985.
  • Breck et al., The ML Test Score: A Rubric for ML Production Readiness, IEEE Big Data 2017.
  • Schelter et al., Automating Large-Scale Data Quality Verification, VLDB 2018.
  • Armbrust et al., Lakehouse: A New Generation of Open Platforms that Unify Data Warehousing and Advanced Analytics, CIDR 2021.

What to learn next

What to learn next

These follow on from what you just read.

  • Data Engineering for AI

    Streaming data

    Streaming means handling data that never stops arriving, where events show up late and out of order, so any number you report is only true so far.

  • Data Engineering for AI

    Building a data pipeline

    A pipeline is every step from raw records to a trained model written down as code, with checks between the steps that stop the run when the data is wrong.

  • Feature and Data Pipelines in Production

    Point-in-time correctness

    Point-in-time correctness means every number shown to a model was actually known at that exact moment in the past, never filled in using something from later.