The Batch, Streaming, and Orchestration Mental Model
Every data pipeline, whatever tool builds it, is answering the same three questions.
Search across all documentation pages
Every data pipeline, whatever tool builds it, is answering the same three questions.
How much data does one run of this pipeline process - a bounded chunk, or an unbounded, ongoing stream?
What decides when and in what order each step of this pipeline runs?
And what happens the second time a step runs on the same input, whether because of a retry, a backfill, or a redelivered message?
Data Engineering Basics shows the concrete API for building batch and streaming jobs; this page is about the three concepts underneath every one of them - batch versus streaming as a trade-off, the DAG as a dependency graph rather than a schedule, and idempotency as the property that makes failure recovery safe.
A batch pipeline processes a bounded, known-size set of data - "yesterday's orders," "this month's log files" - and finishes; running it again later processes the next bounded chunk.
A streaming pipeline processes an unbounded sequence of events as they arrive, with no natural "end" to wait for, which means it has to produce useful output continuously rather than after a full pass over the data.
This is a genuine trade-off, not just two implementations of the same idea: batch can look at all the data before deciding an answer (a full day's total, a complete join), while streaming has to answer with only what's arrived so far, accepting that late-arriving data may need to be reconciled afterward.
A DAG (directed acyclic graph) is the second core idea: it's a graph of tasks and dependencies - "task B needs task A's output" - with no cycles, and critically, a DAG defines what depends on what, not when things run.
A common misreading is to think of a DAG as a fixed sequence, like a numbered list of steps; in reality, a scheduler is free to run any two tasks with no dependency between them in parallel, or in either order, as long as every task's declared dependencies finish first.
Orchestration is the layer that turns a DAG definition into actual execution: it decides which tasks are ready to run (their dependencies are satisfied), retries failed tasks, and tracks each run's state - this is distinct from plain scheduling, which just decides when a job starts, with no notion of dependency between jobs.
Idempotency is the property that running an operation twice with the same input produces the same result as running it once - and it is the property that makes every retry, backfill, and redelivered message safe rather than dangerous.
The batch/streaming choice and the idempotency requirement are more connected than they first appear: any streaming system that promises reliable delivery under failure has to redeliver messages sometimes, and only idempotent processing makes that redelivery harmless instead of duplicating data.
Streaming systems mostly offer at-least-once delivery by default - a message might be processed more than once after a crash and restart, because acknowledging "I processed this" and "commit the result" cannot be made perfectly atomic across two different systems (the message broker and the output sink) without extra coordination.
Exactly-once processing, where it's claimed, is usually achieved by combining at-least-once delivery with an idempotent sink - for example, writing results keyed by a message ID so a duplicate write simply overwrites the same row instead of adding a second one - not by some lower-level guarantee that duplicates never occur.
A watermark is how streaming systems reason about "done enough": since a stream never truly ends, a watermark is a heuristic claim that all data up to some point in event-time has probably arrived, letting windowed aggregations (like "events per 5-minute window") finalize a result while accepting that a small fraction of very late data may be missed or handled separately.
Batch orchestrators face a parallel version of the same idempotency problem: if a nightly job crashes halfway through and is retried, appending its partial output a second time would double-count rows - the standard fix is to make each run's output keyed by a stable partition (a date, a batch ID) and have the job overwrite that partition rather than append blindly.
DAG-based orchestrators (Airflow, Prefect, Dagster) also encode how much history a task run "belongs to" via a data interval - the logical time range a run represents - which is what lets a backfill re-run last month's DAG runs and get last month's data back, rather than accidentally processing today's data under yesterday's label.
# a DAG expresses dependency, not execution order:
# extract_orders and extract_customers have no dependency on each other,
# so an orchestrator is free to run them concurrently
extract_orders >> transform_orders
extract_customers >> transform_orders # transform waits on BOTH upstream tasks
transform_orders >> load_warehouse # load only starts once transform finishesAt scale, the batch-versus-streaming choice stops being binary and becomes a spectrum: micro-batch systems (structured streaming in PySpark, for instance) process small, frequent bounded chunks, trading some of streaming's latency for batch's simpler bounded-data reasoning.
Choosing between Airflow, Prefect, and Dagster is less about which DAG engine is "better" and more about which unit of abstraction matches the team's mental model - Airflow centers on tasks and operators, Prefect centers on plain Python functions with light annotation, and Dagster centers on data assets (the outputs a pipeline produces) rather than the tasks that produce them.
Data validation belongs as a first-class node in the DAG, not an afterthought: a pipeline that loads bad data on schedule is often worse than one that fails loudly, because downstream consumers (dashboards, ML features, other pipelines) may trust the output without knowing it's wrong.
Observability for pipelines has to answer questions batch monitoring and streaming monitoring ask differently: batch asks "did this run finish, and how did its row counts compare to history," while streaming asks "how far behind is this consumer" (consumer lag) and "how many messages are stuck in a dead-letter queue."
File format choice interacts with all of this: columnar formats like Parquet support partition pruning and predicate pushdown that make idempotent partition-overwrite patterns cheap, while row-oriented formats make selective reprocessing of just one partition far more expensive.
| Approach | Strength | Weakness | Best Fit |
|---|---|---|---|
| Batch | Simple reasoning, can see all data before answering, easy backfills | Latency measured in minutes-to-hours; stale by definition between runs | Scheduled reporting, historical backfills, non-time-critical aggregation |
| Micro-batch | Bridges batch simplicity with near-real-time delivery | Still adds a scheduling/trigger latency floor; not truly event-driven | Near-real-time dashboards where seconds-to-minutes latency is acceptable |
| Streaming | Lowest latency, reacts to individual events | Requires watermarks/windowing, harder to reason about "done," state management overhead | Fraud detection, real-time alerting, event-driven architectures |
| DAG orchestration (Airflow/Prefect/Dagster) | Explicit dependency graph, retries, backfills, scheduling in one system | Adds an operational system to run and monitor; overkill for a single linear script | Multi-step pipelines with real task dependencies across sources |
Not directly - a DAG defines dependencies between tasks, and an orchestrator is free to run any two tasks with no dependency between them concurrently or in either order, as long as every task's declared upstream dependencies have finished first.
Scheduling decides when a job starts (a cron-like trigger), while orchestration additionally tracks dependencies between tasks, retries failures, and manages the state of a whole DAG's execution - a scheduler with no dependency model is not an orchestrator.
Failures are routine in distributed pipelines - workers crash, networks partition, messages get redelivered - and idempotency is what turns "this task ran twice" from a data-corrupting event into a harmless no-op, which is what makes automatic retries and backfills safe to run without manual cleanup.
A watermark is a heuristic claim that all events up to a certain point in event-time have probably arrived, which lets a windowed aggregation (like a 5-minute count) finalize and emit a result instead of waiting forever for data that might theoretically still arrive late.
By writing output scoped to a stable partition key - a date, a run ID, a data interval - and overwriting that partition rather than appending to it, so a retried run replaces its own prior (possibly partial) output instead of duplicating rows next to it.
Airflow centers its model on tasks and operators connected in a DAG; Prefect centers on ordinary Python functions with light decoration for dependency tracking; Dagster centers on the data assets a pipeline produces rather than the tasks that produce them - all three still orchestrate DAGs of dependent work underneath.
Micro-batch fits when near-real-time output matters (seconds-to-minutes, not hours) but the team wants to keep reasoning about small bounded chunks rather than building full streaming infrastructure with watermarks and long-lived stateful consumers.
A pipeline that "succeeds" while loading invalid, incomplete, or duplicated data is often more damaging than one that fails loudly, because downstream consumers have no signal to distrust the output - which is why data validation is treated as a DAG node with its own pass/fail state, not a side check.
Consumer lag is the gap between the latest message produced to a stream and the latest message a given consumer has processed; a growing lag means the consumer can't keep up with incoming volume, which is an early warning sign before a pipeline visibly falls behind or drops data.
Only in the trivial sense that a strict sequence is technically a (linear) dependency graph - the practical value of an actual DAG orchestrator shows up once there are genuinely independent branches that could run in parallel, or tasks that need retries, backfills, and dependency-aware scheduling that a plain script has no framework for.
Columnar formats like Parquet support partition pruning, so overwriting or reprocessing just one partition (one date, one key range) touches only the relevant files; row-oriented formats like CSV usually require rewriting or scanning a whole file even to fix one partition, making idempotent partial reprocessing much more expensive.
Stack versions: This page was written for Python 3.14 (stable) / 3.13 (maintenance); the concepts (batch/streaming, DAGs, idempotency) are tool-agnostic and not tied to a specific orchestrator or stream-processing library version.
Reviewed by Chris St. John·Last updated Jul 15, 2026