Summary
Batch and streaming are the two fundamental execution models for data pipelines. Batch processes bounded, finite datasets at scheduled intervals, optimizing for throughput and simplicity. Streaming processes data as an unbounded, continuous flow, optimizing for latency and real-time insight. The choice between them — or a hybrid micro-batch approach — determines pipeline architecture, tooling, cost model, and operational complexity.
Problem
Data arrives continuously from source systems, but downstream consumers have varying freshness requirements. Running everything as streaming adds operational overhead and cost; running everything as batch introduces unacceptable latency for time-sensitive use cases.
- When does a 15-minute data delay become a business problem?
- How do you handle late-arriving data and out-of-order events?
- How do you balance infrastructure cost against freshness SLAs?
Solution
Match the execution model to the latency requirement:
Data Sources
│
├──► [ Batch Pipeline ]
│ Trigger: schedule (hourly, daily)
│ Input: bounded file/table snapshot
│ Output: complete result set
│ Tools: Spark, dbt, SQL, Airflow
│
├──► [ Micro-Batch ]
│ Trigger: fixed short interval (seconds to minutes)
│ Input: small bounded windows
│ Output: incremental result
│ Tools: Spark Structured Streaming, Delta Live Tables
│
└──► [ Streaming Pipeline ]
Trigger: event arrival
Input: unbounded event log
Output: continuously updated result
Tools: Apache Flink, Kafka Streams, Beam
Batch — Load a complete snapshot, apply transformations, write results. Failures are cheap to retry; compute is predictable; joins across large datasets are straightforward.
Micro-batch — Process small time windows on a tight schedule (seconds to minutes). Most of batch’s simplicity, most of streaming’s freshness. Spark Structured Streaming and Delta Live Tables operate this way. The right default for most latency requirements between one minute and fifteen minutes.
Streaming (event-driven) — Process each event or small group as it arrives. Requires explicit handling of windowing, watermarks, late data, and state management. Justified when latency must be sub-minute or when the pipeline itself is the product (e.g., real-time fraud scoring).
Windowing in Streaming
Streaming pipelines must define how to group events into computable units:
Tumbling Window (non-overlapping):
[─── 5m ───][─── 5m ───][─── 5m ───]
Sliding Window (overlapping):
[──── 10m ────]
[──── 10m ────]
[──── 10m ────]
Session Window (gap-based):
[events]─gap─[events]─────gap─────[events]
- Tumbling: Fixed-size, non-overlapping. Aggregation over discrete periods.
- Sliding: Fixed-size, overlapping at a step interval. Moving averages.
- Session: Grouped by inactivity gap. User session analytics, clickstream.
Late Data and Watermarks
Streaming pipelines must define a watermark — the maximum tolerated event-time delay. Events arriving later than the watermark are either dropped, routed to a side-output, or trigger a window correction depending on the engine and business requirement.
# Apache Flink — 10-second watermark tolerance
stream
.assign_timestamps_and_watermarks(
WatermarkStrategy
.for_bounded_out_of_orderness(Duration.of_seconds(10))
.with_timestamp_assigner(lambda event, _: event.event_time)
)
The same tolerance-window idea holds in a TypeScript consumer that has no framework-level watermark support — event time is compared against a rolling maximum, and anything past the tolerance is routed to a side-output instead of the main aggregation:
interface DomainEvent {
eventTime: number; // epoch millis, assigned by the producer
payload: unknown;
}
const LATE_TOLERANCE_MS = 10_000;
class WatermarkGate {
private maxEventTime = 0;
// Returns "onTime" for events within tolerance of the current watermark,
// "late" for events that arrived after their window already closed.
classify(event: DomainEvent): "onTime" | "late" {
this.maxEventTime = Math.max(this.maxEventTime, event.eventTime);
const watermark = this.maxEventTime - LATE_TOLERANCE_MS;
return event.eventTime < watermark ? "late" : "onTime";
}
}
function partitionByWatermark(events: DomainEvent[]) {
const gate = new WatermarkGate();
const onTime: DomainEvent[] = [];
const late: DomainEvent[] = [];
for (const event of events) {
(gate.classify(event) === "onTime" ? onTime : late).push(event);
}
return { onTime, late };
}
WatermarkGate tracks the highest event time seen so far as a proxy for “now” in event time, rather than trusting wall-clock arrival order — the same principle Flink’s WatermarkStrategy applies internally, made explicit for pipelines built on plain Kafka consumers.
Live Playground
Experiment with the pattern below. The two tabs show a fixed-interval batch job (before.ts) that groups events strictly by wall-clock arrival — a late event silently lands in whichever batch happens to be open when it shows up — and a watermark-aware equivalent (after.ts) that classifies events against event time and routes late arrivals to a side-output instead of corrupting the wrong window.
When to Use
Batch:
- Reporting and analytics where hourly or daily freshness is acceptable
- Large-scale joins across full historical datasets
- ETL pipelines feeding data warehouses on a known schedule
- Workloads where idempotent reruns and simple failure recovery are required
Micro-batch:
- Dashboards and operational metrics with one- to fifteen-minute SLAs
- Incremental ingestion into lakehouses (Delta Lake, Iceberg)
- When streaming complexity is unjustified but pure batch latency is too high
Streaming:
- Fraud detection and anomaly alerting where seconds matter
- Real-time personalization, recommendations, and bidding systems
- Event-driven architectures where downstream systems react to events immediately
- IoT telemetry processing at the edge or near-edge
Avoid streaming when:
- The business question is answered daily or hourly — streaming adds cost and complexity without value
- The team lacks operational experience with stateful stream processing
- Late data handling requirements are undefined — silent correctness bugs are hard to detect
Trade-offs
| Benefit | Cost |
|---|---|
| Batch: simple, high-throughput, easy reruns | Minutes-to-hours latency; stale data |
| Batch: full dataset available for joins and aggregations | Peak compute demand; resource spikes at schedule time |
| Streaming: sub-second to sub-minute latency | Stateful processing, watermarks, and late-data logic required |
| Streaming: continuous output; no scheduled delays | Operational complexity; harder to debug and test |
| Micro-batch: latency/simplicity balance | Window boundaries introduce artificial delay; still not true event-time |
Related Patterns
- Medallion Architecture — Defines the layered zones that both batch and streaming pipelines write into
- Pure Functions — Transformation logic should be pure regardless of execution model, enabling testability and reruns
- Schema-Driven Validation — Applied at ingestion in both models to catch malformed events early