Data Pipelines
ETL and ELT, idempotency, incremental extraction and change data capture, late data, backfills, orchestration, data contracts and the quality checks that catch silent corruption.
A data pipeline is a repeatable process that moves data from where it is produced to where it is used, transforming it on the way. The useful definition is narrower than the diagram suggests: a pipeline is a set of units of work with declared inputs, declared outputs, a trigger, and a defined behaviour when it is run twice. Everything hard about pipelines — correctness after a failure, backfills, late data, silent corruption — is a consequence of how carefully those four things were specified.
ETL, ELT and why the order changed
In ETL, data is extracted, transformed by a separate processing tier, and loaded into the target already shaped. In ELT, raw data is loaded first and transformed inside the target system using its own query engine. The shift toward ELT happened because storage became cheap enough to keep raw data and warehouse engines became powerful enough to do the transformation work, not because one is intrinsically better.
The practical difference is reproducibility. In ELT the untransformed input is still in the system, so a logic bug can be fixed by rerunning a transformation over data you still have. In ETL, if the transform was wrong, the input may be gone and the only recovery is re-extracting from a source that has since changed. That single property — can I rebuild this table from something I still possess? — decides more pipeline designs than any performance consideration.
The four properties every task needs
Idempotency
A task is idempotent when running it twice produces the same result as running it once. This is not a nicety; retries are guaranteed to happen, from schedulers, from operators and from partial failures, so a non-idempotent task is a corruption waiting for a bad night. The mechanics are well known: make writes replace a partition rather than append to a table, use a merge keyed on a natural or surrogate identifier rather than a blind insert, and write to a temporary location before swapping it in atomically so a crashed job leaves no half-written output. If a task cannot be made idempotent, it must at least be made detectable — a run identifier stamped on every row so duplicates can be found later.
Incremental extraction
Full refreshes are simple, obviously correct, and stop being viable as data grows. Incremental extraction replaces them and introduces the pipeline's most common correctness bug. A high-watermark approach — select rows where the updated timestamp is greater than the last run's maximum — fails whenever the source can write a row with a timestamp earlier than the watermark: clock skew between source nodes, long transactions committing after the read, or a timezone-naive column compared against a value in another zone. The mitigations are to overlap the window deliberately (re-read a safety margin and rely on idempotent merges), to use a monotonic sequence rather than wall-clock time where one exists, and to reconcile periodically against a full count.
Change data capture avoids the problem entirely by reading the source database's transaction log instead of querying it, which also captures deletes — the events a timestamp query can never see, because a deleted row has no updated timestamp. The cost is a tighter operational coupling to the source: log retention, schema changes, and the initial snapshot all become your concern.
Late and out-of-order data
Data arrives after the period it belongs to. A mobile client was offline; a partner delivers yesterday's file at noon; a correction is issued for last quarter. A pipeline must have an explicit answer to this, and there are only three honest ones: reprocess affected periods on a rolling window, accept the record into the current period and keep the event time as a column so consumers can decide, or reject it and record the rejection. The dishonest answer is to silently drop anything outside the current window, which produces numbers that quietly disagree with the source.
Backfill as a first-class path
Every pipeline will need to recompute history: a bug is found, a definition changes, a new column is added. If backfill is a different code path from the scheduled run, it will be wrong when you need it most. The design that avoids this is parameterising every run by the period it processes, so that "run yesterday" and "run every day of last March" are the same code with different parameters. Backfills should also be throttled — an unbounded backfill competing with production for the same compute turns a data fix into an availability incident.
Orchestration
An orchestrator manages dependencies, scheduling, retries and visibility. It should coordinate work, not perform it: pulling data through the scheduler process makes a coordination service into a single-node compute engine with no recovery semantics and no scaling story.
The important distinction is between time-triggered and data-triggered execution. Running at a fixed hour assumes the input is ready by then; it is the default and it is the source of a whole category of silent-wrong-answer failures, because a job that runs on an incomplete input usually succeeds. Triggering on the arrival or completeness of the input — a sensor, an event, a published marker — is strictly more correct, and worth the extra machinery for anything whose output people act on. Where time triggers stay, add an explicit freshness assertion so an early run fails loudly instead of computing on partial data.
Data contracts, testing and observability
Pipelines break in two ways: they fail, or they succeed with wrong data. The first is easy; the orchestrator tells you. The second is the one that reaches a dashboard, and only explicit checks catch it.
Test at three levels. Schema: required columns exist with expected types, and unexpected new columns are surfaced rather than absorbed. Row-level: primary keys are unique and non-null, foreign keys resolve, enumerations contain only known values, numeric ranges are sane. Aggregate: row counts and key business totals fall within an expected band relative to recent history, and freshness is within its target. These are cheap to implement as assertions that run as part of the pipeline and fail the run, which is far better than a monitor that notices afterwards.
A data contract makes the schema and its compatibility rules an agreement between the producing team and the consuming pipeline, rather than something the consumer reverse-engineers from yesterday's payload. Contracts are what allow an upstream team to move fast without breaking you: they can add fields freely, and must negotiate removals and type changes.
Lineage — knowing which outputs derive from which inputs — turns "this number looks wrong" from an archaeology project into a lookup, and turns "we must delete this user's data" from a guess into a list.
Common failure modes
- Non-idempotent appends. A retried task doubles a day's rows; nothing errors, and the totals are simply wrong.
- Deletes that never propagate. A timestamp-based incremental load cannot see deletions, so the target accumulates rows the source no longer has.
- Timezone drift. A naive timestamp compared against a value in another zone shifts every boundary by hours, producing a plausible but wrong zero or double-count at the edges of each period.
- Success on empty. The source returns nothing because of an auth failure or an outage, the pipeline loads zero rows and reports success, and the dashboard shows a flat line that looks like a business event.
- Silent schema absorption. A column changes type and the loader coerces it, so a numeric field becomes text and every downstream aggregate quietly stops matching.
- Poison records with no dead-letter path. One malformed row fails the batch forever; the fix is a quarantine location plus an alert, not a wider try/except.
- The pipeline nobody owns. It runs, it is green, and nobody can say what consumes its output or whether it still needs to exist.
When not to build a pipeline
Copying data is a liability: the copy can be stale, it can disagree, and it must be maintained. Before building a pipeline, check whether a view over the source, a query federation, or a read replica answers the question — the cheapest pipeline is the one that does not exist. Build one when the read pattern would harm the source system, when data from several systems must be joined, when history must be preserved that the source overwrites, or when a transformation is expensive enough that it should be computed once and reused.
When you do build, the ordering that keeps you out of trouble is: land raw first, make every task idempotent, parameterise by period so backfills are ordinary, and assert on the data before anyone sees it. See big data architecture for where pipelines fit in the wider platform, data warehousing for their most common destination, and streaming for the continuous variant.