Detail chart · Data Pipeline

Batch vs Streaming

Pattern ◆◆◆◇◇

Choose between bounded batch processing and unbounded stream processing based on latency, throughput, and consistency requirements.

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.

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]

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:

Micro-batch:

Streaming:

Avoid streaming when:

Trade-offs

BenefitCost
Batch: simple, high-throughput, easy rerunsMinutes-to-hours latency; stale data
Batch: full dataset available for joins and aggregationsPeak compute demand; resource spikes at schedule time
Streaming: sub-second to sub-minute latencyStateful processing, watermarks, and late-data logic required
Streaming: continuous output; no scheduled delaysOperational complexity; harder to debug and test
Micro-batch: latency/simplicity balanceWindow boundaries introduce artificial delay; still not true event-time