Why processing only what changed makes pipelines faster, cheaper, and easier to scale.
You have a table with one billion records. Every day, one million new ones arrive. Would you really reprocess the entire billion just to handle the new million?
That is the problem incremental processing solves.
It sounds obvious when you say it out loud. Yet a surprising number of production pipelines still work the naive way: read everything, transform everything, write everything, every single run. It works at first. Then the data grows, the bill grows, the runtime grows, and one day the nightly job is still running when people open their laptops in the morning.
This article explains what incremental processing really is, the details that make it work in practice, and when the old fashioned full refresh is still the right answer. No hype. Just the pattern, the mechanics, and the judgment calls.
Incremental processing rests on one simple idea:
Do only the work that is necessary.
If only one million rows changed, your pipeline should touch roughly one million rows. Not one billion. The savings are not small. Processing one percent of the data usually means the job runs in a small fraction of the time, costs a small fraction of the compute, and finishes long before anyone notices it ran at all.
But here is the honest part. Incremental processing is not free. You trade compute for carefulness. A full refresh is simple because it has no memory: it recomputes everything from scratch, so it cannot get out of sync. An incremental pipeline has to remember what it already processed, detect what changed since then, and merge those changes into the existing result without duplicating or losing anything.
That memory, that change detection, and that merging are where the real engineering lives. Let us walk through each one.
Before you can process only the changes, you have to know what changed. There are three common ways, and every incremental pipeline you will ever meet is built on one of them.
A timestamp column. The most common approach. Every row carries an updated_at timestamp, and the pipeline asks: give me everything where updated_at is greater than the last value I saw. Simple, cheap, and good enough for most operational tables. The catch is that it depends on the source system reliably maintaining that column. If someone backfills rows without updating the timestamp, those rows become invisible.
A sequence or ID. Some systems give every change a monotonically increasing number: an auto-increment ID, a log sequence number, an offset. The pipeline stores the highest number it has processed and asks for everything above it. This is more reliable than a timestamp because sequences do not have clock skew, timezone bugs, or "forgot to update the column" problems. If your source offers one, prefer it.
Change Data Capture (CDC). The most complete approach. Instead of asking the table what changed, you read the changes directly from the system's own change log. Delta Lake has this built in as Change Data Feed: enable it once, and every insert, update, and delete is recorded with its own metadata, which you can read with table_changes(). We walk through exactly how this works, with runnable code, in our lesson on incremental processing and Change Data Feed. CDC also captures deletes, which timestamps and sequences simply cannot see.
A quick way to choose: if the table only ever receives appends, a timestamp or sequence is enough. If rows can be updated or deleted, you want CDC, or you accept that your incremental pipeline quietly misses those changes.
Detecting changes is only half the job. The other half is applying them.
Suppose an order changes status from "pending" to "shipped" overnight. If your pipeline just appends the new row, your table now has two rows for the same order, and every downstream report double counts it. This is the single most common bug in hand-rolled incremental pipelines.
The answer is merge logic, often called an upsert: if the row already exists, update it; if it does not, insert it. In Delta Lake this is a first-class operation:
MERGE INTO orders AS target
USING incoming_changes AS source
ON target.order_id = source.order_id
WHEN MATCHED AND source.updated_at > target.updated_at THEN
UPDATE SET *
WHEN NOT MATCHED THEN
INSERT *Notice the extra condition: only update when the incoming row is actually newer. That one clause protects you from a whole class of out-of-order delivery problems. We go deeper on MERGE patterns, including deletes and slowly changing dimensions, in the SCD patterns lesson and the Delta Lake lesson.
Real data does not arrive on time. A mobile app syncs hours later. A partner system sends a correction file for yesterday. A sensor comes back online and uploads three days of backlog.
If your pipeline processes "everything since the last run" exactly once, late data falls through the cracks. The standard fix is an overlap window: instead of starting exactly at your last checkpoint, you deliberately re-read a little history. If your watermark says you processed up to midnight, you reprocess from, say, 48 hours before midnight every run.
Reprocessing the same rows twice is safe precisely because of the merge logic from step two. MERGE is idempotent: running it twice with the same input produces the same result as running it once. That is the beautiful combination. The overlap window catches the stragglers, and the idempotent merge makes the overlap harmless.
Jobs fail. Networks drop, clusters get preempted, someone deploys a bug. The question is never whether your pipeline will fail mid-run. The question is what happens when it retries.
A well-built incremental pipeline has two properties:
Correct checkpoints. The pipeline records how far it has processed, and it only advances that checkpoint after the work is durably written. Advance the checkpoint too early and a failure silently skips data. Advance it too late and the retry reprocesses data, which is fine as long as your writes are idempotent. When in doubt, checkpoint late.
Idempotent writes. Every write can be safely repeated. MERGE gives you this for row-level changes. For partition-level work, Delta Lake's replaceWhere overwrite gives you the same guarantee at the partition level: re-running the backfill for November 15th replaces November 15th, exactly, no matter how many times it runs.
If you take one sentence from this article to an interview or a design review, take this one: an incremental pipeline is only as trustworthy as its checkpoints are conservative and its writes are repeatable.
Incremental processing fits naturally into the medallion architecture, but it behaves a little differently at each layer.
Bronze is almost always append-mostly. Raw data lands, and you add new files or new rows. Timestamps or sequences work well here, and if you are on Delta Lake, Change Data Feed or structured streaming with availableNow = True makes this nearly automatic.
Silver is where merge logic earns its keep. Cleaned, deduplicated, current-state tables need upserts. This is where the MERGE pattern above lives, and where data quality checks belong, since a bad incremental update can quietly corrupt a table that everything downstream trusts. Our data quality lesson covers how to put those guardrails in place.
Gold is where the judgment calls happen. A small aggregate table over a huge fact table is often cheaper to rebuild entirely than to maintain incrementally. A large fact table feeding many dashboards almost always wants incremental updates. The layer does not decide for you; the size of the data and the cost of recomputation do.

Incremental processing is the default for large, frequently updated tables. It is not a religion. A full refresh is still the right answer when:
The mature pattern is hybrid: incremental for speed, periodic full refresh for correctness. You get the daily efficiency and a built-in repair mechanism.
https://youtu.be/RTQsZwfm3Ns
Incremental processing is not an optimization you bolt on at the end. It is a design decision with four moving parts: how you detect changes, how you merge them without duplicating, how you catch late arrivals with an overlap window, and how checkpoints plus idempotent writes make failures boring instead of dangerous.
Get those four right, and the one-billion-row table stops being scary. The nightly job finishes in minutes, the compute bill shrinks, and the pipeline scales with the size of your changes instead of the size of your history.
Do only the work that is necessary. It is a good rule for pipelines, and honestly, not a bad rule for everything else too.
If you are new here, the Start Here lesson is the best first step into the book. Everything above is easier once the foundations are calm and clear.