Data Pipeline
Data Pipeline
Definition: An automated sequence of steps that moves data from one or more sources to one or more destinations, typically involving ingestion, validation, transformation, storage, and serving. It is the umbrella concept for this entire corner of data engineering — every other tool in this folder is a specialized piece that slots into one stage of this same larger shape: Apache Airflow orchestrates it, Apache Kafka moves data through it continuously, Apache Spark transforms it at scale, and Change Data Capture (CDC) feeds it incremental source changes. A pipeline can run in scheduled batches, continuously as a stream, or as a hybrid of both, but the underlying mental model — data flowing through a sequence of well-defined, coordinated stages — stays the same regardless of which specific technology implements each one.
Architecture
- Ingestion — pulling data from source systems (production databases, third-party APIs, event streams, file drops) into the pipeline, either as a bulk extract or an incremental capture via something like Change Data Capture (CDC)
- Validation — checking incoming data against expected schema, types, and business rules before it’s allowed further downstream, the core job of Data Quality and Validation
- Transformation — cleaning, joining, aggregating, and reshaping raw data into the form downstream consumers actually need, commonly done at scale by an engine like Apache Spark
- Storage — landing transformed (and often raw) data in a durable destination, usually a lake or warehouse (see Data Lake vs Data Warehouse) or a managed platform like Snowflake
- Serving — the final stage, exposing data to whatever consumes it: BI dashboards, ML feature stores, downstream applications, or another pipeline entirely
- Orchestration — a coordinator such as Apache Airflow sits above the stages, sequencing them, enforcing dependencies, retrying failures, and giving engineers one place to see the whole pipeline’s health
- Idempotency — every stage should be safe to re-run against the same input without producing duplicate or corrupted output, since retries are a routine operating condition, not a rare edge case (see Idempotency)
- Observability — logging, metrics, and alerting attached to every individual stage, not just a single pipeline-level “succeeded or failed” flag, so failures and quality issues surface immediately (see Observability and Monitoring)
The canonical shape, with an orchestrator coordinating every stage from above:
Why It Matters
- Almost every analytics dashboard, ML model, and operational report depends on a pipeline running correctly and on schedule behind the scenes, invisible right up until it breaks
- It’s the connective tissue for this entire vault section — Apache Kafka, Apache Spark, Apache Airflow, Change Data Capture (CDC), and ETL vs ELT are all implementations of one or more pipeline stages, not competing or unrelated ideas
- Getting the shape right — which stages exist, how decoupled they are, where state lives — determines how easy the system is to debug, extend, and recover from failure for years afterward
- A well-designed pipeline degrades gracefully: a failure in one stage should be visible and recoverable, not a silent, cascading outage
- The batch-versus-streaming decision made at the pipeline level (see Batch vs Stream Processing) ripples through every downstream choice, from storage format to how consumers are allowed to query the data
Under the Hood: Idempotency and Failure Recovery
Pipelines fail constantly, in ways that usually have nothing to do with the transformation logic being wrong — a network blip mid-transfer, a source database timing out, a worker node getting evicted mid-job. The design question is never “how do we prevent failure,” which is impossible to guarantee at scale over enough runs, but “how does a stage recover from failure without corrupting or duplicating data.” That property is Idempotency: re-running a stage on data it has already processed should land in the same end state as running it exactly once.
A few mechanisms turn idempotent recovery from a theory into something operable:
- Checkpointing — a stage periodically records how far it has gotten (for example, “orders up to ID 48213 are loaded”), so a restart resumes from the last checkpoint instead of from the very beginning
- Watermarking — in streaming contexts specifically, a watermark tracks the point in event-time up to which the pipeline considers data complete, letting it decide when a window of results is safe to finalize versus still waiting on late-arriving events
- Idempotent writes — upserts keyed on a stable ID, rather than blind appends, so replaying the same batch twice overwrites existing rows instead of duplicating them
- Dead-letter queues — records that repeatedly fail validation or transformation get routed aside instead of blocking or crashing the entire run, so one bad row doesn’t take down a whole day’s data
Without these, a retry after a partial failure either reprocesses everything from scratch — wasteful, and often slow enough to blow past an SLA — or reprocesses nothing and silently drops whatever was in flight when the failure hit. Both outcomes are worse than the original failure.
The same five conceptual stages take a very different physical shape depending on whether the pipeline is batch or streaming:
Batch treats the pipeline as a job that starts, finishes, and stops; streaming treats it as a process that never stops running at all — see Batch vs Stream Processing for the deeper trade-offs behind choosing one shape over the other.
Comparison: Batch Pipeline vs Streaming Pipeline vs Micro-Batch Pipeline
| Batch Pipeline | Streaming Pipeline | Micro-Batch Pipeline | |
|---|---|---|---|
| Latency | Minutes to hours, often once a day | Sub-second to a few seconds | Roughly seconds to a couple minutes |
| Complexity | Lower — straightforward to reason about and test | Higher — ordering, state, and late-arriving data all need explicit handling | Medium — batches the streaming problem into small, manageable chunks |
| Typical orchestration | Apache Airflow or a plain cron schedule | The stream processor’s own runtime, e.g. Apache Kafka consumers | Apache Spark Structured Streaming in micro-batch mode |
| Typical failure recovery | Re-run the whole job, or the failed task, from a checkpoint | Resume from a committed offset, reprocessing only what wasn’t acknowledged | Re-run just the failed micro-batch, bounded by design to a small time window |
Common Pitfalls
- No idempotency, so a retry after a partial failure duplicates every record that had already been written before the failure hit
- Tightly coupling stages so a failure in transformation cascades backward into ingestion or forward into storage, taking down an entire run instead of just one piece
- No monitoring beyond “did the job exit with a success code,” so silent data quality issues — nulls, duplicates, schema drift — go unnoticed for days or weeks
- Designing only for the happy path and never testing what happens when a source is late, empty, or malformed, then discovering the gap for the first time in production
- Treating schema changes at the source as someone else’s problem, so an upstream team renaming a field silently breaks everything downstream
- Keeping transformation logic in a person’s head or a wiki page instead of versioned code, so nobody can safely change or debug the pipeline later
- Reprocessing historical data by hand under pressure during an incident instead of designing backfills as a repeatable, first-class operation from the start
Worked Example
Trace a single nightly run of a pipeline that processes yesterday’s e-commerce orders, stage by stage:
| Stage | What happens | Example detail |
|---|---|---|
| Trigger | Orchestrator starts the DAG once upstream data is expected to be ready | Airflow schedule fires at 0 2 * * * |
| Ingestion | Extract yesterday’s rows from the orders and customers tables | Roughly 40,000 new order rows pulled via incremental extract |
| Validation | Check for nulls in required fields, valid foreign keys, and a plausible row-count range | 12 rows rejected for a missing customer ID, routed to a dead-letter table |
| Transformation | Join orders to customers, compute order totals, tag each order with a region | A Spark job reshapes raw rows into a clean, analytics-ready table |
| Storage | Write the transformed table into the warehouse, partitioned by date | Loaded into a Snowflake table partitioned on order_date |
| Serving | Downstream dashboard queries refresh automatically against the new partition | Analysts see yesterday’s orders when they check the dashboard each morning |
Real-World Use
- Analytics reporting pipelines — nightly or hourly jobs that populate the dashboards executives and analysts check every morning
- ML feature pipelines — computing and materializing features on a schedule or continuously, so a model sees consistent, freshly computed inputs at both training and inference time
- Real-time personalization pipelines — streaming pipelines that update a user’s recommendations or search ranking within seconds of a new interaction
- Data sync between systems — keeping a CRM, a warehouse, and a search index consistent with each other, often via Change Data Capture (CDC) feeding a shared pipeline
- Regulatory and financial reporting — pipelines that must produce auditable, reproducible numbers on a fixed schedule, where idempotency and lineage matter as much as speed
Best Practices
- Design for idempotency from the start — retrofitting safe re-runs onto a pipeline already in production is far harder than building it in from day one
- Monitor data quality at every stage, not just overall pipeline “success,” since a job can exit cleanly while still producing wrong or incomplete data (see Data Quality and Validation)
- Decouple stages so a failure in one is isolated and recoverable rather than cascading through the entire run
- Version pipeline logic alongside the data it produces, so any output can be traced back to the exact code that generated it
- Make backfills and reprocessing a designed, repeatable capability, not a one-off script written under pressure mid-incident
- Prefer boring, well-understood tools that fit a stage’s actual requirements over defaulting to whatever’s most powerful or fashionable
FAQ
Is a data pipeline the same thing as ETL? Not exactly — ETL (and ELT) describe one common pattern for moving and transforming data, while “pipeline” is the broader umbrella term for any automated, staged movement of data, batch or streaming, ETL-shaped or not; see ETL vs ELT.
Does every pipeline need an orchestrator like Airflow? No — a single, simple scheduled job can just be a cron entry. Orchestrators earn their keep once there are multiple interdependent tasks, retries, and cross-task failure handling to manage.
What’s the difference between a pipeline and a workflow? The terms overlap heavily in practice; “pipeline” tends to emphasize the flow of data through stages, while “workflow” or DAG emphasizes the dependency graph of tasks moving that data — Airflow orchestrates workflows that implement pipelines.
Can one pipeline be both batch and streaming? Yes — many production architectures run a streaming path for low-latency needs alongside a batch path over the same sources for heavier historical reprocessing and correction.
Why do pipelines break so often? Mostly because they sit at the boundary between systems nobody on the data team fully controls — upstream schemas change, APIs go down, traffic spikes — so failures usually come from external drift, not bugs in the pipeline’s own logic.
How is a pipeline different from just writing a script? A script typically runs once, by hand, and a human sorts things out if it breaks; a pipeline is expected to run unattended, on a schedule or continuously, and needs its own retry, monitoring, and recovery behavior built in to survive that.
History
- ETL’s roots trace to 1970s-80s enterprise data warehousing, when nightly batch jobs first moved data out of transactional systems and into separate reporting databases
- Through the 1990s and 2000s, hand-rolled scripts and cron jobs were the default way to chain these steps together, with little shared tooling across companies
- The “big data” era of the late 2000s and early 2010s, led by Hadoop and MapReduce, pushed pipelines toward distributed processing over commodity clusters as data volume outgrew single-machine ETL jobs
- Purpose-built orchestrators emerged through the 2010s — Apache Airflow, open-sourced by Airbnb in 2014, became the de facto standard for defining pipelines as code instead of tangled cron chains
- Streaming platforms like Apache Kafka, originally built at LinkedIn around 2011, shifted part of the industry from “pipeline as a scheduled batch job” toward “pipeline as a continuous, always-on process”
- The current era leans toward lakehouse-oriented storage and streaming-first design, blurring the old hard line between batch ETL and real-time processing rather than treating them as separate systems
Common Interview Questions
- “Walk me through how you’d design a pipeline to process daily order data” — expect each stage named in order (ingest, validate, transform, store, serve) with tool choices justified at each one
- “How would you make a pipeline stage idempotent?” — expect discussion of upserts keyed on a stable ID, checkpointing, and avoiding blind appends
- “How do you handle a schema change from an upstream source?” — expect mention of schema validation, versioning, and failing loudly rather than silently ingesting bad data
- “When would you choose streaming over batch for a given pipeline?” — expect the answer to hinge on actual latency requirements weighed against streaming’s added operational complexity
- “How would you debug a pipeline that’s producing wrong numbers but hasn’t technically failed?” — expect an answer centered on per-stage data quality monitoring, not just job-level success or failure status
Related Terms
- Apache Airflow
- Apache Kafka
- Apache Spark
- ETL vs ELT
- Change Data Capture (CDC)
- Batch vs Stream Processing
- Data Quality and Validation
Example
A nightly pipeline pulls yesterday’s orders from a production database, joins them with customer data, validates the result against expected row counts and null checks, and loads the output into a warehouse table analysts query each morning. Underneath, an orchestrator like Airflow triggers each step in sequence, retries the extraction if the source database briefly times out, and pages the on-call engineer if validation fails outright. Swap the nightly database pull for a Kafka topic and the same five stages — ingest, validate, transform, store, serve — become a streaming pipeline instead, no different in shape, just continuous instead of scheduled.
Referenced by