A full reload read 1,000,050 rows to pick up 200 changes. A watermark alone read 199 and missed a late order; a watermark with a 1-hour lookback read 1,588, wrote exactly 200 and matched the source.
python3 --version.pip install "duckdb>=1.0"git clone https://github.com/DayanEbrar0X/data-anatomy.ai.git cd data-anatomy.ai
python3 -m venv .venv source .venv/bin/activate
pip install -r requirements.txt # or just this lesson: pip install "duckdb>=1.0"
cd data-engineering/19-incremental-loads python3 src/load.py
from orders import night_one, day_two, report
def full(db, _):
n = db.execute("CREATE OR REPLACE TABLE wh "
"AS FROM src").fetchone()[0]
return n, n # read all, write all
def incremental(db, minutes):
read = db.execute("""CREATE TABLE new AS
FROM src WHERE updated_at > (SELECT wm
FROM state) - INTERVAL (?) MINUTE""",
[minutes]).fetchone()[0]
wrote = db.execute("""MERGE INTO wh
USING new USING (id) WHEN MATCHED
AND new.updated_at > wh.updated_at
THEN UPDATE WHEN NOT MATCHED THEN INSERT
""").fetchone()[0]
db.sql("UPDATE state SET wm = "
"(SELECT max(updated_at) FROM wh)")
return read, wrote
for name, load, minutes in [
("full", full, 0),
("watermark", incremental, 0),You reload a million rows to pick up 200 changes. Incremental loads fix that. Read only what changed. Like a bookmark: you start where you stopped.
Our source holds a million orders, each with an updated_at time. Last night we loaded them all, and saved the newest time: midnight. The high-watermark. Today, 200 rows changed.
One arrived late. A phone that was offline syncs an order stamped 11:40 last night. Older than the watermark, but it landed after the load. The full reload replaces the warehouse table with the whole source.
It reads every row, and writes every row. The incremental load copies rows updated after the stored watermark, minus a lookback window, in minutes, for late data. Then MERGE into the warehouse on the order id. A match is updated only if the incoming row is newer.
Anything unmatched is inserted. Then move the watermark up, to the newest row. Three runs: full, watermark only, and a one hour lookback. Each starts from last night, gets today's 200 changes, and checks the result.
Let's run it. Full: 1,000,050 read, 1,000,050 written. Match. Watermark only: 199 read and written.
No match. The late order was missed. Lookback: 1,588 read, only 200 written. Match.
MERGE skipped the old rows it re-read. Now the work grows with your changes, not your table. Size the lookback to your latest data, and always MERGE, so re-reads are harmless. One gotcha: watermarks never see deletes, so reload in full now and then.
Save the mark, look back a little, merge on the key. Not a million rows for 200 changes.
Read the lesson on GitHub →