ETL Pipeline: The Extract That Fails Silently Is Worse Than the One That Crashes
Incremental extraction, idempotent loads, late-arriving data, backfills, and the data quality checks that decide whether anyone trusts the warehouse.

An ETL pipeline that crashes is a solved problem. It alerts, someone restarts it, the data catches up. The pipelines that cause real damage are the ones that report success while loading 90% of the rows, or the right number of rows from the wrong columns, or yesterday’s rows twice.
That asymmetry drives the whole design. The engineering effort belongs on making silent failures impossible to detect-by-accident and on making the loud failures cheap to recover from.
The scale, and where the difficulty is
50,000 source tables
400 pipelines
2,000 tables in the warehouse
daily volume
5 billion rows
~4 TB
incremental: 400 GB changed per run
schedule
hourly for 300 pipelines
daily for 100
the pipeline that is actually hard
not the 4 TB
it's the one reading a source with no
reliable ordering, late updates, and a
schema that changed without notice
The scale numbers are the easy part; every architecture handles them. What makes ETL hard is that sources are inconsistent in ways you do not control, and the pipeline is the only place where those inconsistencies become visible.
Extract: full, incremental, and the watermark problem
full extract
read the whole table every run
4 TB, 30 minutes, most of it unchanged
simple, correct, and unaffordable
-> use for small tables, reference data,
anything under ~1M rows
incremental by primary key
WHERE id > :last_max_id
only if ids are monotonic and never reused
works when ids are a sequence
breaks when ids come from a sequence that
was reset by a restore
incremental by updated_at
WHERE updated_at > :watermark
the standard answer, and the subtle one
change data capture
read the write-ahead log or replication stream
every change, in order, no polling
the correct answer when the source supports
it, at the cost of operational complexity
Incremental extraction on updated_at is what most pipelines use and it has exactly one serious problem: the timestamp is not the same as the time the row became visible.
source row updated at 10:00:00
transaction commits at 10:00:05
pipeline runs at 10:00:30 with watermark 10:00:00
-> the row is not visible, watermark advances
-> 10:00:00 < 10:00:30
-> this row is now never extracted
-> a record silently missing, forever
The fix is a safety window, and the size of it is a real parameter with a real trade-off:
watermark = now() - safety_window
safety_window = 1 hour
-> a row that becomes visible up to 1h late
is still captured
-> the pipeline re-reads the last hour
on every run (wasted work, bounded)
safety_window = 0
-> data loss, silent and permanent
safety_window = 7 days
-> safe, and now every run re-reads a week
And even with a safety window, some data arrives later than the window. Deletes and hard updates are worse: a row modified twice where the second modification predates your last watermark is missed, and you have no way to detect it from the source. The honest position is that updated_at extraction can miss updates and the only real fixes are CDC, a full reconciliation run, or a source that gives you a change log you can trust.
Which is why the reconciliation run exists. Periodically — nightly, weekly, monthly — re-read a source fully and compare against what the pipeline has, then repair. It is expensive, it is unglamorous, and it is the only thing that catches extraction that silently lost data. Treat it as a scheduled correctness guarantee rather than as an occasional project.
Transform: where correctness is decided
The transform is a pure function in the useful sense: same input, same output. Making it genuinely pure is worth real effort because it makes every downstream property cheap.
Make the logic explicit and testable. A transformation in a language you can unit test, with fixtures of awkward source data, catches problems that a runtime error never will.
-- the kind of thing that silently drops rows
SELECT
customer_id,
-- inner join: customers with no orders vanish
o.order_id
FROM orders o
INNER JOIN customers c ON c.id = o.customer_id
-- LEFT JOIN keeps them, with NULLs
-- which is correct depends on intent
Inner versus left join is the single most common source of quietly missing rows in a warehouse, and it is invisible because the pipeline succeeds and the counts look plausible. Join semantics deserve explicit review on every transformation where a relationship is optional.
Type and precision changes matter. Money in a float, an integer column that overflows into NULL, a timestamp that shifts timezone on the way in. These are the failures that pass every count-based check.
Do the transformation in the engine, not in application code, when the volume justifies it. Moving the same logic from a Python process to a warehouse query changes the failure modes: the database applies one consistent snapshot, and you stop depending on whether a source row was modified between two of your own queries.
Load: idempotency is not optional
append
INSERT INTO events SELECT * FROM stage
- no key
- a retry appends the same rows again
- double counts forever
- never use for anything countable
merge on key
MERGE INTO events USING stage
ON events.event_id = stage.event_id
WHEN MATCHED THEN UPDATE SET ...
WHEN NOT MATCHED THEN INSERT ...
- idempotent
- a retry produces the same result
- the standard answer
delete and insert partition
DELETE FROM events WHERE partition_date = '2025-10-16'
INSERT INTO events SELECT * FROM stage WHERE ...
- idempotent
- fast on columnar storage
- the answer for append-only fact tables
truncate and reload
TRUNCATE events; INSERT ...
- idempotent
- brief window of emptiness
- wrong for anything serving queries
Choose one per table and be consistent, because mixing them across a single table produces counts that depend on the load order.
The staging pattern makes this safe and is worth the extra storage:
source -> stage (raw, append whatever)
|
| validate BEFORE touching target
| - row counts vs expectation
| - null rates on required columns
| - key uniqueness in stage
| - referential sanity
v
stage -> transform -> target
|
| merge / delete-insert
v
target -> published
Validating the staged data before the merge is what makes a bad run harmless. A pipeline that merges first and validates after has already corrupted the target, and “the check failed” is now a data incident rather than a rejected run.
A subtlety worth naming: a merge is idempotent for the same input, but if the pipeline re-extracts with a different safety window it may see a newer version of a row, and merging is correct. What is not idempotent is an append of a source snapshot that has already been applied. Both look like “re-running the pipeline” from the operator’s side and they behave differently.
Handling late, late-arriving data
A row that arrives today about an event from three days ago is normal, and handling it is a design decision, not an accident.
event_time 2025-10-13
loaded_at 2025-10-16
options
1. ignore it
the daily aggregate for 10-13 is now
wrong and you do not know
2. append to the fact table, leave the
aggregate alone
the fact table is correct, the rollup
is stale
3. re-partition: update the affected day
correct, and you are re-writing a
partition that queries may have read
4. separate correction path
facts are append-only, corrections live
in their own table and a view subtracts
the prior version
most systems need 1 or 2
a finance system reporting on a closed
period needs 3 or 4
The decision is a function of whether a published number is allowed to change. If dashboards are read once and discarded, option 1 is fine and cheap. If there is a regulatory or contractual close, you need corrections, and the way to do that without rewriting history is a correction table with the version and an effective timestamp, which a view uses to resolve the current value.
Whatever the option, the pipeline must know which one it applies and report it, because “some data arrived late and I did nothing about it” should be a visible, deliberate state rather than a discovery during a reconciliation.
Data quality checks that mean something
the checks, by usefulness
strong
- expected vs actual row count, per table,
per run, with a threshold
- uniqueness on the primary key
- referential: every order_id in fact_orders
exists in dim_customers
- null rate on required columns
- freshness: max(loaded_at) within an SLA
- value ranges and distributions
(revenue cannot be negative; a customer's
age cannot go backwards)
weak
- "the job succeeded"
- "the table is not empty"
- "row count is greater than yesterday"
(a bug that halves every load passes this)
where they run
- on stage, before merge: reject the run
- on target, after merge: alert
- both, because they catch different bugs
The primary key uniqueness check is the highest-value single check in the whole discipline, and it should exist on every table. It catches duplicate appends, double-counting from a re-run, and merge logic that is not doing what was intended.
Distribution checks catch the rest. A revenue column that is 40% negative might be a legitimate refund, or it might be a sign error in the transform, and no count check will tell you. Comparing today’s distribution against a rolling baseline catches it in a way a fixed threshold does not.
Quarantine beats reject for anything you cannot fully validate. Rather than failing the whole run because 3 rows out of 5 billion are bad, route them to a quarantine table with the reason, publish the rest, and surface the count. The failure mode to avoid is a pipeline that either halts on one bad row or silently drops it, and quarantine is the option that is neither.
Backfill is a feature, not a rerun
A backfill re-runs history: a schema change last month, a bug found in a transform, a new column computed retroactively.
why it is dangerous
- it runs the same job as the live pipeline
- it competes for the same warehouse capacity
- it may write to partitions users are
actively querying
- one failure can leave partitions in
a half-updated state
what makes it safe
- isolated: a separate job, its own
connection and resource pool
- bounded: explicit partition list and
date range, never "everything"
- rate limited: throttle to leave capacity
for the live pipeline
- ordered: newest partition last, so a
partial backfill leaves the freshest
data correct
- idempotent: same merge logic, so a
rerun of a backfill is safe
- observable: its own metrics, its own
alerting, visible separately from the
live pipeline so nobody mistakes a
backfill spike for an incident
- reversible: keep the pre-backfill
partition, so it can be restored
The reversibility point is the one teams skip and the one that matters. A backfill that overwrites correct historical partitions with a wrong result is a data incident that takes weeks to unpick, and “we can recompute from source” is not always true. Snapshotting the target partitions before the backfill costs storage and removes an entire category of incident.
Orchestration and the DAG
extract_orders ──┐
extract_customers┤
extract_products ├──> transform_facts ──> load_facts ──> quality_facts
│ │ │
extract_payments┘ v v
dim_customer publish
|
freshness checks on each extract ─────────────┘
The DAG’s job is dependencies and recovery, and two properties are where the design lives.
Idempotent tasks with retries. A task that is retried must produce the same result. Every task therefore needs a natural key and a merge rather than an append, and the DAG needs a retry policy per task type with a cap and a dead-letter path.
A clear definition of a run’s completeness. Is the pipeline done when every task succeeded, or when the data is fresh? These are different and conflating them produces a “successfully completed” run with missing partitions. Mark the run complete only when the freshness check passes on the target tables, not when the last task returned zero.
Airflow and similar orchestrators are the usual choice and the usual source of the second problem: a DAG with 400 tasks, a webserver that serialises them, and a scheduler that silently stops catching up. The specific failure to watch is scheduled-run overlap — a run that takes longer than its interval starting a second copy of itself, with two runs writing the same partitions. Prevent it with a run identity and a lock, and detect it, because it produces duplicated data that looks like correct data for a while.
Failure stories worth testing
Move the source clock forward and back
Confirm the watermark behaviour. Clock adjustment on the source is the most common cause of silently missed rows, and it looks identical to “no data changed”.
Run a pipeline twice concurrently
Confirm the second one is rejected or isolated. Overlapping runs write the same partitions and produce duplicates that satisfy every count check.
Retry a job after a partial load
Confirm the target is unchanged or correctly converged. This is the single most important idempotency test and it should exist for every pipeline.
Load 90% of yesterday’s rows
Confirm the row-count check catches it. If the check compares to yesterday and yesterday was also 90%, it passes — check against the source, not against history.
Truncate a required column upstream
Confirm the null-rate check rejects the run rather than loading 5 billion nulls.
Change a column’s meaning without changing its name
Confirm nothing catches it. This is the failure quality checks cannot detect, and it is why distribution checks and semantic ownership exist.
Re-run a backfill over a live month
Confirm the live pipeline’s partitions are untouched and the freshest partition is written last. Partial backfills that break the live path are the standard backfill incident.
Corrupt a partition with a wrong backfill and try to restore
Confirm the pre-backfill snapshot exists. If it does not, the data is gone and the only recovery is a full re-extract, if the source still has it.
Flood the warehouse from a backfill and a live pipeline
Confirm the backfill is rate limited and the live pipeline keeps to its freshness SLA. Backfills that starve the live path are why backfills need their own resource class.
Put one bad row in a 5-billion-row batch
Confirm the row is quarantined with a reason and the run publishes. Neither halting the whole pipeline nor silently dropping the row is acceptable.
Kill the scheduler and let the DAG fall behind
Confirm catch-up works and does not run concurrently with the next scheduled run.
A production-ready architecture
50,000 source tables
|
| full for small, incremental with a
| safety window for large, CDC where
| the source supports it
v
+---------------------------+
| stage (raw, append) | nothing trusted yet
+------------+--------------+
|
v
+---------------------------+
| VALIDATE BEFORE MERGE | row count vs source
| - count vs expectation | key uniqueness
| - key uniqueness | null rate on required
| - null rate | referential sanity
| - value ranges | value distributions
| -> quarantine bad rows |
| with a reason |
+------------+--------------+
| pass
v
+---------------------------+
| transform | pure functions, tested,
| in the engine, one | explicit join semantics,
| consistent snapshot | no float for money
+------------+--------------+
|
v
+---------------------------+
| load (idempotent) | merge on key, or
| | delete+insert partition
| NEVER append to a |
| countable table |
+------------+--------------+
|
v
+---------------------------+
| target + published view | freshness SLA on
| | max(loaded_at)
+---------------------------+
orchestrator DAG
- task-level retries with a cap
- run identity + lock, no overlapping runs
- "run complete" = freshness passes,
not "last task returned"
- backfills: separate resource class,
rate limited, newest partition last,
pre-backfill snapshot, own metrics
late data
- ignore | append | re-partition | corrections
- chosen deliberately per table, reported
A sensible delivery checklist:
- Pick one load strategy per table and use it consistently. Merge on key for dimension-like tables, delete-and-insert partition for append-only facts.
- Validate on the stage before the merge, so a bad run is rejected rather than corrupting the target.
- Run a primary key uniqueness check on every table. It is the highest-value single check in the discipline.
- Check row counts against the source, not against yesterday, and be suspicious of any check that compares to history.
- Use a safety window on every
updated_atextraction, and document what the window costs. - Schedule periodic full reconciliation runs. Incremental extraction can miss updates and this is the only thing that catches it.
- Make the transform a tested, pure function with explicit join semantics. Review every inner join where the relationship is optional.
- Quarantine invalid rows with a reason and publish the rest. Do not halt the pipeline and do not drop silently.
- Lock runs by identity so two runs cannot write the same partitions, and alert on any overlap.
- Define run completion as a freshness check on the target, not as the last task returning successfully.
- Treat backfills as a separate resource class with rate limits, ordering, a pre-change snapshot, and their own metrics.
- Decide the late-arriving-data policy per table deliberately, and report when a row was late and ignored.
Common mistakes
| Mistake | What actually happens | Better decision |
|---|---|---|
| Append to a fact table | Retry double-counts, permanently | Merge on key, or delete+insert partition |
| Validate after the merge | Corruption happens before the check | Validate on stage, reject the run |
updated_at with no safety window |
Rows that commit late are lost forever | Safety window, sized and documented |
| No reconciliation run | Missed updates are never discovered | Scheduled full compare and repair |
| Row count vs yesterday | A bug that halves the load passes | Compare against the source |
| Job succeeded = pipeline healthy | A 90% load looks successful | Freshness check on the target |
| Inner join on an optional relationship | Rows silently disappear | Explicit join choice, reviewed |
| Float for money | Rounding errors, unreproducible sums | Decimal, fixed precision |
| No key uniqueness check | Duplicates accumulate unnoticed | Check every table, every run |
| Halting on one bad row | One row blocks 5 billion good ones | Quarantine with a reason, publish the rest |
| Silently dropping bad rows | Unknown data loss, discovered in an audit | Quarantine table, count reported |
| No run lock | Two runs, same partitions, duplicates | Run identity + lock, alert on overlap |
| Backfill on the live job | Backfill starves the live pipeline | Separate resource class, rate limited |
| Backfill with no snapshot | A wrong backfill is unrecoverable | Snapshot target partitions first |
| Backfill in arbitrary order | Partial backfill leaves freshest data wrong | Newest partition last |
| “Done” when tasks return | Run marked complete with missing partitions | Complete on freshness, not on return |
| Unbounded task retries | Broken task loops forever, queue stalls | Retry cap, dead-letter, visible |
| Late data ignored silently | A published number is wrong, nobody knows | Deliberate policy, reported |
| No distribution checks | A sign error passes every count check | Compare distributions to a baseline |
| Meaning changed, name did not | Nothing detects it | Semantic ownership, reviewed definitions |
| One pipeline for 50,000 tables | A single failure takes out everything | Partition by domain, isolate failures |
| Transform in application code | Inconsistent snapshots, hard to test | In the engine, testable pure functions |
The complete story in one minute
An ETL pipeline is a distributed system with a database at each end, and the failures that hurt are the ones that report success. A pipeline that loads 90% of yesterday’s rows, or the right number of rows with a sign error in one column, is worse than one that crashes, because a crash is a five-minute fix and a silent wrong number is a quarter of distrust.
Idempotency is the foundation. Pick one strategy per table — merge on a primary key, delete-and-insert the partition, or truncate-and-reload — and use it consistently, because an append with no key double-counts the first time a job retries. Validate on the staging table before the merge, so a bad run is rejected rather than corrupting the target, and quarantine invalid rows with a reason instead of either halting five billion good rows or dropping them without trace.
Extraction is subtler than it looks. updated_at incremental extraction with no overlap can lose a row whose transaction becomes visible after the watermark advanced, and a clock adjustment on the source produces the same symptom. An overlap plus an idempotent merge reduces that race; it does not catch a correction older than the overlap or a hard delete. A source log position/CDC stream is a stronger cursor, and independent reconciliation is still required because connectors, retention, and destination logic can fail too.
And the checks have to be real. Row counts compared against the source rather than against yesterday, primary key uniqueness on every table, null rates on required columns, and distribution checks that catch a negative revenue column that a count check cannot. Green means fresh, which is not the same as correct, and a freshness signal alone will let a pipeline that truncates pass every check.
Backfills get their own resource class, their own metrics, a pre-change snapshot, newest-partition-last ordering, and rate limiting. And a run is complete when freshness passes on the target, not when the last task returned zero.
extract with a safety window, reconcile periodically
stage -> validate -> merge (idempotent, key-based)
complete means fresh, not "returned successfully"
backfills are isolated, ordered, snapshot, rate limited
The hard part was never moving the bytes. It was building the checks that let you trust a green pipeline without opening the data.


