Batch vs Stream Processing

Batch vs Stream Processing

Definition: Two fundamental models for processing data: batch processing accumulates data over a period and processes it as a large, finite chunk on a schedule, while stream processing handles each event individually and continuously as it arrives. The core tradeoff is latency versus completeness — batch waits until it can see a full, stable set of data before computing anything, while stream processing trades away that completeness for results measured in seconds or milliseconds. A third option, micro-batching, sits between the two, grouping events into very small time-boxed batches rather than committing fully to either extreme. Neither model is strictly “better” — the right choice depends on how fresh an answer needs to be, and how much operational complexity a team can absorb to get there.

How Each Model Works

  • Batch — accumulate, then process: data lands in storage and simply sits there until a scheduled job sweeps it up and processes the whole set in one run
  • Batch — bounded input: each run operates over a known, finite dataset — “yesterday’s rows,” “this hour’s file” — so it can compute exact, stable aggregates without worrying about more data arriving mid-computation
  • Batch — scheduling: jobs are commonly triggered by a scheduler like Apache Airflow, running nightly, hourly, or on some other fixed cadence, often as one stage in a larger Data Pipeline
  • Batch — cheap, deterministic reprocessing: because the input is a stable, addressable dataset, rerunning a job with corrected logic or backfilled data is comparatively simple — “just run it again” produces a trustworthy answer
  • Stream — unbounded input: data is modeled as an endless sequence of events with no defined end, so a streaming job never waits for “all of the data” because there’s no such thing
  • Stream — per-event or micro-batch processing: each event is processed individually within milliseconds of arriving, or grouped into very small micro-batches — Spark Structured Streaming being the canonical example — trading a few seconds of latency for higher throughput and simpler failure recovery
  • Stream — windowing: since summing “forever” isn’t meaningful, streaming systems group events into windows — tumbling, sliding, or session-based — to compute rolling metrics like “orders in the last 5 minutes”
  • Stream — watermarks: a watermark tracks how far behind event time a pipeline is willing to wait for late data, letting an otherwise-endless window eventually close and emit a result
  • Both — different optimization targets: batch systems are tuned to maximize throughput over a known volume, while streaming systems are tuned to minimize per-event latency, and pushing either one to behave like the other usually means a different architecture, not just a config change

Why It Matters

  • Freshness is a product decision, not just a technical one: whether users see “as of last night” or “as of right now” data often determines whether a feature is viable at all — a fraud check that only runs nightly isn’t really a fraud check
  • The choice sets the complexity floor for the whole system: streaming pipelines commonly require durable logs, stateful processing, and careful handling of failure and replay — infrastructure a pure batch pipeline never has to think about
  • It drives cost, in both directions: batch jobs can run on cheap, ephemeral compute that spins up once a day; streaming jobs commonly run always-on infrastructure provisioned for peak load around the clock
  • It determines how errors get found and fixed: a wrong number in a batch report is usually just rerun; a wrong number already served from a live streaming dashboard may need an explicit correction or retraction
  • It shapes team structure and skills: streaming systems commonly demand deeper operational expertise — monitoring consumer lag, tuning windows, reasoning about exactly-once semantics — that batch-only teams may not have built up
  • It underlies most “real-time” product claims: anything marketed as live, instant, or real-time is, underneath, a claim about which of these two models — or which hybrid — is actually running

Under the Hood: Windowing and Late-Arriving Data

Streaming systems face a problem batch systems mostly avoid: data never stops arriving, so an aggregation has to decide, on its own, when it’s allowed to stop waiting and actually emit a result. A handful of ideas make this tractable:

  1. Tumbling windows — fixed-size, non-overlapping time intervals (e.g. every 5 minutes); each event belongs to exactly one window
  2. Sliding windows — fixed-size but overlapping intervals (e.g. a 5-minute window recalculated every 1 minute); a single event can contribute to several windows at once
  3. Session windows — variable-length windows defined by a gap in activity (e.g. close the window after 30 minutes of silence from a given user), commonly used for things like tracking one browsing session
  4. Event time vs. processing time — event time is when something actually happened; processing time is when the pipeline observed it; network delays, retries, and outages routinely make the two diverge, sometimes by minutes or hours
  5. Watermarks — a running heuristic for “how far behind event time are we still willing to wait,” e.g. “close the 2:00-2:05 window once we’ve seen event time 2:05 plus a 2-minute grace period”
  6. Late data handling — events that arrive after their window’s watermark has already passed are dropped, routed to a separate late-data output, or trigger a retraction and recomputation of an already-emitted result, depending on how the engine is configured

Real production systems very rarely commit to one pure model — this is where hybrid architectures come in, blending a slow, accurate path with a fast, approximate one.

Lambda architecture, proposed by Nathan Marz, runs a batch layer and a speed layer against the same raw data source at once: the batch layer periodically recomputes accurate, complete views over all historical data, the speed layer computes fast, approximate views over only the most recent slice, and a serving layer merges both — queries hit the fast view for anything the latest batch run hasn’t caught up to yet. Kappa architecture simplifies this by dropping the separate batch layer entirely: everything, including full historical reprocessing, runs through the same stream-processing engine, with “batch” reprocessing implemented as simply replaying old events from a durable log — commonly Apache Kafka — back through that same streaming job.

Comparison: Batch vs Stream vs Micro-Batch

DimensionBatchMicro-BatchStream (Event-at-a-Time)
LatencyMinutes to hours, or longerSeconds to low minutesMilliseconds to low seconds
ThroughputVery high — optimized for large, uniform sweepsHigh — amortizes overhead across small groupsStrong per-engine, though per-event overhead is higher
ComplexityLowest — simple to write, test, and rerunModerate — some streaming concerns, simpler failure modelHighest — state, windowing, out-of-order handling
Fault recoverySimple — rerun the whole jobModerate — replay a small batchComplex — checkpointing, exactly-once semantics
Data completeness at query timeComplete, as of the last runNearly complete, seconds staleComplete as of now, individual results may later need correction
Typical toolsApache Spark batch, Hadoop MapReduce, Airflow-scheduled jobsSpark Structured StreamingApache Kafka Streams, Apache Flink, Apache Storm
Typical use caseNightly reports, ML training, historical backfillsNear-real-time dashboards, rolling aggregatesFraud detection, alerting, live monitoring

Common Pitfalls

  • Choosing streaming when a nightly batch job would have met the actual latency requirement — wiring up Kafka, a stream processor, and stateful windowing logic for a report nobody checks more than once a day is pure overhead
  • Ignoring late and out-of-order data until it causes a visibly wrong number in production — treating event time and processing time as interchangeable works fine until a delayed event silently skews an aggregate that’s already been emitted
  • Letting streaming job state grow unbounded — a windowed aggregation or join that never expires old keys will eventually exhaust memory, especially with high-cardinality keys like user or session IDs
  • Treating micro-batch as “real-time” when the requirement genuinely needs sub-second, event-at-a-time processing — a 30-second micro-batch window is not the same guarantee as true streaming, and conflating the two causes missed SLAs
  • Underestimating how much harder streaming is to debug and reprocess compared to batch — there’s no simple “just rerun yesterday” when the pipeline has no natural start or end
  • Skipping idempotency in stream consumers — at-least-once delivery is the common default, meaning a consumer will occasionally see the same event twice, and a non-idempotent write silently turns that into duplicated data
  • Building a Lambda architecture’s dual code paths without keeping batch and streaming logic in sync — when the two implementations drift apart, batch and real-time views quietly disagree and often nobody notices until someone compares them directly
  • Over-provisioning streaming infrastructure for peak load that rarely materializes — always-on clusters sized for a worst-case spike can cost far more, year-round, than the occasional burst they exist to absorb

Worked Example: Two Requirements, Two Implementations

Requirement 1 — “Give the finance team total sales per region, every morning.”

  1. Orders land in a raw events table throughout the day, untouched
  2. A batch job, scheduled via Apache Airflow for 1 a.m., reads the full previous day’s rows
  3. The job groups by region, sums order totals, and writes one row per region into a reporting table
  4. Finance opens a dashboard at 8 a.m. reading yesterday’s fully-settled numbers — refunds and cancellations from late in the day are already reflected

Requirement 2 — “Block a card the moment it looks like fraud.”

  1. Each card swipe publishes an event to an Apache Kafka topic the instant it happens
  2. A stream processing job consumes the topic continuously, maintaining a rolling window of each card’s recent activity, e.g. transaction count in the last 5 minutes
  3. When a new event pushes that rolling count past a threshold, the job emits an alert within milliseconds of the triggering swipe
  4. A downstream service consumes the alert and declines the next authorization request — waiting for a nightly batch job here would mean the fraud already happened hours before anyone found out

The same underlying question — “how much activity just happened” — gets answered by a full, stable table scan in one case and a continuously-updated rolling window in the other, purely because the acceptable latency differs by roughly six orders of magnitude.

Real-World Use

  • Nightly business reporting — end-of-day sales totals, inventory reconciliation, and financial close processes almost always run as batch jobs, since “as of midnight” is a perfectly acceptable freshness bar
  • ML model training — training a model over months of historical data is a textbook batch workload, even when the resulting model is later served in real time
  • Fraud and anomaly detection — banks and payment processors stream every transaction through rules or models that must flag suspicious activity within milliseconds, not after the fact
  • Real-time dashboards and monitoring — operational dashboards showing live order counts, server error rates, or active users depend on stream processing to stay meaningfully “live”
  • Alerting pipelines — infrastructure monitoring tools stream metrics and logs continuously so an on-call engineer is paged within seconds of a threshold breach, not the next morning
  • Change data capture feeds — Change Data Capture (CDC) tools stream database row-level changes into downstream systems continuously, keeping replicas and search indexes in near-real-time sync

Best Practices

  • Default to batch unless a specific, articulated latency requirement genuinely needs streaming — batch is simpler to build, test, monitor, and reprocess, and most “we need real-time” requests actually mean “within the hour,” which batch can still deliver
  • Design streaming consumers to be idempotent from day one — assume at-least-once delivery and make repeated processing of the same event a no-op, rather than retrofitting deduplication after duplicates cause damage
  • Decide explicitly how late data will be handled before writing the first windowing rule — dropping it, emitting a correction, and routing it to a side output are all valid choices, but only if the choice is deliberate
  • Consider a hybrid approach, Lambda- or Kappa-style, before committing fully to pure streaming — most “real-time” requirements only apply to a narrow, recent slice of data, with everything older served perfectly well by batch
  • Monitor consumer lag, not just throughput — a streaming job that’s technically “up” but falling further behind the event source every hour is failing its actual purpose even if no error is ever thrown
  • Keep windowing and watermark logic simple and well-documented — these are the parts of a streaming system most likely to produce subtly wrong numbers that pass code review but fail on real-world edge cases

FAQ

Is stream processing always faster than batch? Per-event latency, yes — but batch systems are often dramatically higher-throughput for the same total volume, since they avoid the per-event bookkeeping overhead streaming requires.

Can a single system do both batch and stream processing? Yes — engines like Apache Spark and Apache Flink support both models on a shared execution engine, which is a large part of why the strict batch/stream boundary has blurred in recent years.

What’s the difference between micro-batch and true streaming? Micro-batch groups events into very small time-boxed batches, commonly sub-minute, and processes each group at once; true event-at-a-time streaming processes each event as it individually arrives, with no grouping delay at all.

Why not just always use streaming, since it sounds more “real-time”? Because streaming systems are meaningfully more complex to build, operate, and debug — that complexity isn’t free, and a lot of “real-time” requirements don’t actually need sub-second freshness once you ask precisely how fresh is needed.

What happens if a streaming job goes down for an hour? Well-designed systems buffer events durably upstream, commonly in Apache Kafka or another Message Queue, so the job can resume from where it left off and catch up, rather than losing the gap entirely.

Does streaming replace the need for a data warehouse? No — streaming and batch outputs both commonly land in the same warehouse or lake; streaming changes how fresh the data is when it lands, not where it ultimately lives, see Data Lake vs Data Warehouse.

Is streaming always more expensive than batch? Usually, yes, per unit of data processed — always-on compute and the operational overhead of state management typically cost more than a job that runs once and shuts down, though this gap narrows as managed streaming services take on more of that burden.

History

  • 1950s-2000s, batch as the default: early data processing — punch cards, mainframe batch jobs, nightly ETL runs — was batch by necessity, since hardware and storage costs made anything else impractical
  • 2004, MapReduce: Google’s MapReduce paper formalized large-scale batch processing over distributed, commodity hardware, and the open-source Hadoop implementation made that model broadly accessible, defining the “batch era” of big data
  • Late 2000s-2010s, streaming engines emerge: Apache Storm (2011) and Apache Kafka — originally built at LinkedIn, open-sourced around 2011 — brought practical, distributed stream processing and durable event logs to a much wider audience
  • 2011, Lambda architecture proposed: Nathan Marz proposed running batch and speed layers side by side against the same data, formalizing a hybrid pattern many systems had already been approximating ad hoc
  • Around 2015, the Dataflow model: Google’s Dataflow paper and the resulting Apache Beam project proposed a single unified programming model for bounded and unbounded data alike, influencing how later engines expose windowing and watermarks as first-class API concepts
  • Mid-2010s, Kappa architecture proposed: Jay Kreps, a Kafka co-creator, proposed Kappa as a simplification — treat everything as a stream, and implement “batch” reprocessing as simply replaying the stream from the start
  • Mid-to-late 2010s, unification: Apache Spark’s Structured Streaming and engines like Apache Flink increasingly blurred the batch/stream distinction, offering one programming model and one engine capable of expressing either
  • 2020s, streaming as a default expectation: cloud-managed streaming services and rising real-time product expectations pushed more teams to consider streaming earlier in a system’s design, even in cases where batch would still have sufficed

Common Interview Questions

  • “When would you choose batch processing over stream processing, and vice versa?” — expect a latency-vs-complexity tradeoff answer, anchored to a concrete business requirement rather than a default preference for either
  • “How would you handle a late-arriving event in a streaming pipeline?” — expect a discussion of watermarks, grace periods, and the choice between dropping, correcting, or side-outputting late data
  • “What’s the difference between event time and processing time?” — expect a clear distinction plus an example of when they diverge, like a mobile app that buffers events offline and uploads them hours later
  • “Explain the Lambda architecture and why someone might move away from it” — expect the dual batch/speed-layer description, plus the maintenance cost of keeping two codebases in sync as the usual reason teams migrate toward Kappa or a unified engine
  • “How does a streaming system provide exactly-once processing guarantees?” — expect mention of idempotent writes and transactional offsets or checkpointing, plus the caveat that “exactly-once” is usually really “effectively-once”
  • “What does unbounded state growth mean in a streaming job, and how do you prevent it?” — expect mention of expiring old window state, bounding key cardinality, and setting explicit time-to-live on stored aggregates

Example

A batch job recalculates monthly sales totals every night, scanning the full prior day’s orders table and writing one settled number finance can trust; a streaming pipeline instead updates a live “orders in the last hour” dashboard as each order happens, trading that same certainty for a number that’s usually right within a few seconds. Both jobs can read from the exact same underlying order events — the only difference is whether the result is computed once over a finished, bounded set, or continuously over a set that never stops growing.

Dig deeper