Stream vs. Batch: Processing Data at Scale

Event streams vs batch jobs, lambda vs kappa architectures, and how to choose the right one when the data never stops arriving.

Advanced · 18 min read

Why this matters

Every system you've built so far answers questions about the present. Sooner or later somebody asks about the past — "how many orders shipped last Tuesday?" — and about right now — "is the checkout flow broken this second?" One question wants a settled, final answer. The other wants an answer while the data is still arriving. Batch processing and stream processing are the two families of tools for those two questions, and picking the wrong one is how you end up with a fraud report that's twelve hours late, or a "real-time" dashboard that's really a cron job with ambition.

Video: Beyond Batch Processing: Streaming to Accelerate Value for Data | Amazon Web Services — Amazon Web Services
AWS explains the real value of processing data continuously instead of waiting on batches, with stateless/stateful processing and windowing basics.

The core distinction: the newspaper stack vs the live ticker

This lesson's running analogy has two halves. The first is a stack of newspapers delivered to your door every morning: finite, complete, countable. You can read it in any order, re-read page one, and know for certain when you've finished it. The second is the live news ticker at the bottom of a 24-hour news channel: it never ends, you can't see tomorrow's headlines, and if you blink, you missed one.

That is the whole distinction:

Batch processing reads the stack. Stream processing watches the ticker. Everything else in this lesson — the tools, the architectures, the trade-offs — falls out of that one difference.

Video: Batch vs Realtime Stream Processing - A Deep Dive — The Geek Narrator
Deep-dive interview contrasting batch and realtime processing on latency, throughput, cost, idempotency, and schema evolution.

Batch: reading the morning stack

Batch processing is the older sibling, and its recipe hasn't changed in twenty years: collect everything into a pile, run one big job over all of it on a schedule, write the results somewhere queryable. The pile is usually a data lake or warehouse; the schedule is hourly, nightly, or weekly; the job is MapReduce or, these days, Spark.

The canonical example: at 2 a.m., a Spark job reads the entire day's worth of order events, joins them against the product catalog, and writes out yesterday's revenue report. Nobody expects the answer at 2:01 a.m. — the job might take an hour — but when it finishes, the answer is complete and correct, computed over every single record.

In newspaper terms: you wait for the full stack to be delivered, sit down with a red pen, and read the whole thing before anyone is allowed to ask you a question. High throughput (crunching terabytes in one job), high latency (answers are hours old by design).

flowchart LR
    Apps["Your services<br/>emit logs and events"]:::client --> Lake[("Data lake<br/>the whole stack, piled up")]:::data
    Lake --> Job["Batch job<br/>runs on a schedule"]:::service
    Job --> Views[("Batch views<br/>precomputed answers")]:::data
    Views --> Dash["Dashboards<br/>answer yesterday's questions"]:::client

Batch shines where correctness and completeness beat freshness: billing, financial reporting, training machine-learning models over months of history, and — critically — reprocessing. If you discover the revenue logic had a bug, you fix the job and re-run it over the whole stack. The stack is still there. That's a superpower the ticker doesn't have for free.

Video: Data Engineering Deep Dive: Batch Processing & Immutable Artifacts — deferstech
Walks through MapReduce internals, HDFS, Spark dataflow engines, and why immutable batch outputs make rollback safe.

Streams: watching the ticker

Stream processing flips the recipe: instead of waiting for the pile, you process each event as it arrives. The events flow through a durable event log — Kafka, which you met in the Message Queues lesson — and a stream processor like Flink reacts to each one, keeps running state, and emits results continuously. Latency drops from hours to seconds or less.

The catch is that the ticker forces you to answer questions the stack never asked. Three of them are worth internalizing:

Event time vs processing time. Every headline on the ticker has two timestamps: the time printed on the story (when the click happened — event time) and the moment it scrolled past your screen (when your system saw it — processing time). They're different because phones go offline, networks stall, and retries happen. Count by processing time and a click from midnight lands in the morning's bucket. Count by event time and the buckets are right — but you have to wait for stragglers, because the 11:59 p.m. headline might scroll past at 12:04 a.m.

Windowing. You can't total "all time" on a ticker that never ends, so you slice it into windows — like clipping the ticker into segments and summarizing each one. Tumbling windows are fixed, non-overlapping blocks (every 5 minutes). Sliding windows overlap (the last 10 minutes, recomputed every minute). Session windows group bursts of activity separated by gaps. Picking the window is picking the question.

Watermarks. Late headlines are a fact of life, so stream processors use watermarks: a marker that says "we're now confident we've seen every event up to 9:00." Once the watermark passes 9:00, the 8:55–9:00 window closes and its result is emitted — without waiting forever for a headline that might never arrive. A watermark is a heuristic, not a promise: a straggler arriving after the watermark is either dropped or routed to a late-data path. Set the grace too tight and you lose data; too loose and your "real-time" answers are just batch with extra steps.

A concrete pass through all three ideas: you're counting ad clicks in 5-minute tumbling windows. The 9:00–9:05 window should close at 9:05 — but a click stamped 9:04:58 (event time) only reaches your processor at 9:05:03 (processing time) because the user's phone was in a tunnel. With a 2-minute watermark, the window stays open until 9:07, the late click lands in the right bucket, and the count is correct. A click stamped 8:59 that arrives at 9:08 misses the watermark and goes to the late-data path. Every streaming number you've ever trusted was produced by exactly this machinery.

flowchart LR
    Events["Clicks, orders,<br/>sensor readings"]:::client --> Log[("Event log<br/>Kafka")]:::data
    Log --> Engine["Stream processor<br/>Flink"]:::service
    Engine --> Live[("Live views<br/>updated continuously")]:::data
    Live --> Alerts["Dashboards and alerts"]:::client

Stream processing also has to survive crashes mid-ticker without double-counting or losing state — Flink does this with periodic checkpoints of its state, so a failed worker restarts from the last checkpoint and replays the log from there. Conceptually: if you doze off during the broadcast, you rewind to where you remember and keep watching. That log-replay trick will matter a lot in a moment.

Video: Chapter 11: Stream Processing Explained — deferstech
The DDIA chapter on stream processing, distilled: Kafka's log-first design, CDC, event time vs processing time, windows, stream joins, and exactly-once — the missing manual for the live ticker.

Lambda: the stack and the ticker, stapled together

In the early 2010s, Nathan Marz — then at Twitter — looked at the two approaches and decided he wanted both: the completeness of batch and the freshness of streams. The result was the Lambda architecture, described in his and James Warren's book Big Data (2015). It has three layers:

flowchart TD
    Src["Events pour in"]:::client --> Log[("Immutable event log<br/>the master dataset")]:::data
    Log --> Batch["Batch layer<br/>recomputes everything, nightly"]:::service
    Log --> Speed["Speed layer<br/>approximates, in real time"]:::service
    Batch --> Serve["Serving layer<br/>merges both answers"]:::service
    Speed --> Serve
    Serve --> App["Your query"]:::client
    Batch -. "same logic,<br/>written twice" .-> Speed

Back to the analogy: you keep the entire newspaper archive and watch the ticker, then staple the two answers together whenever someone asks. It works. It also costs you something real: two implementations of the same logic — a batch job and a streaming job that must agree with each other forever. They drift. Somebody fixes a bug in one and not the other. And when an answer looks wrong, you get to play detective: which layer lied? Lambda buys you both freshness and completeness, and the price is operating two systems that do the same job differently.

Video: Lambda vs Kappa vs Data Streaming Platform — Which One Wins? | Modern Data Engineering Pipelines — Confluent Developer
Explains why Lambda architecture exists: batch layer for correctness plus speed layer for freshness, then its dual-codebase costs.

Kappa: burn the stack, keep the ticker

In 2014, Jay Kreps — one of Kafka's creators — published an essay questioning whether the second system was ever worth it. His proposal, the Kappa architecture, is disarmingly simple: everything is a stream. One event log, one stream-processing engine, no batch layer at all.

The trick is the log-replay you saw earlier. If the log is immutable and retained long enough, you can rebuild any view by replaying history through new code. Bug in the revenue logic? Deploy the fix and rewind — reprocess the whole log with the corrected processor. Need a brand-new view nobody asked for before? Run a new stream job over the old log. The batch layer's superpower (reprocessing) turns out to live in the log, not in the batch job.

flowchart TD
    Src["Events pour in"]:::client --> Log[("Immutable event log<br/>retained for replay")]:::data
    Log --> Job["One stream processor"]:::service
    Job --> Views[("Serving views")]:::data
    Views --> App["Your query"]:::client
    Job -. "new logic?<br/>rewind and replay" .-> Log

Kappa wins on simplicity: one codebase, one system to operate, no "which layer lied" debates. But it has a hard requirement that Lambda doesn't: the log must be kept long enough to replay. Retention becomes a first-class design decision — keep 7 days of log and you can only ever rebuild 7 days of history. And for true bulk computation — scanning years of data to train a model — a streaming replay is a slow, awkward way to do what a batch job does naturally. Kappa replaces the batch layer only when your log retention covers your reprocessing horizon.

Interactive diagram: VsToggle (loads in the app)

Video: Big Data Kappa Architecture Explained — Thomas Henson
Shows how Kappa drops the batch layer entirely so Spark, Flink, and search all consume one streaming layer.

One problem, two solutions: counting ad clicks

Make it concrete. An advertiser asks: "how many clicks did my ad get in the last hour?"

The batch answer: every hour, a Spark job scans the click log, groups by ad, and writes the counts. Exact to the click — and an hour stale. The advertiser adjusting their budget at 3 p.m. is steering by the 2 p.m. numbers. For money, that's fine: invoices get settled by the batch job, where every click is counted exactly once.

The stream answer: a Flink job counts clicks in 5-minute tumbling windows on event time, with a 2-minute watermark. The advertiser watches a dashboard update every few seconds and pauses a burning campaign before lunch. A click arriving 10 minutes late misses its window and lands in a late-data stream — the dashboard is approximately right in real time, and the nightly batch reconciliation makes the invoice exactly right.

Notice what just happened: the real production answer was a hybrid — streams for operations, batch for money. That's the decision guide in action, and it's how most serious ad platforms actually work. Freshness where it changes behavior, exactness where it changes balances.

Video: Turning the database inside out with Apache Samza — Martin Kleppmann (Strange Loop)
The DDIA author himself: treat data as a stream of immutable facts, and a live counter becomes just a rollup of the click log — the core mental model behind real-time aggregation.

Delivery semantics: the fine print on the ticker

The ticker raises one more question the stack never asks: when a headline scrolls past twice, do you count it twice? Network retries, crashed workers, and replays all mean stream processors routinely see the same event more than once. So every stream system picks a delivery semantic:

This is the same guarantees conversation from the Message Queues lesson, now wearing a compute hat. The practical rule carries over unchanged: assume at-least-once, make your logic idempotent, and save "exactly-once" for the narrow paths where duplicates are genuinely expensive.

Video: External systems & data streaming | The Duchess & The Doctor #3 ft. Anna McDonald & Matthias J. Sax — Confluent Developer
Kafka experts dig into exactly-once versus retries, stateful snapshots, and correctness trade-offs in event-driven systems.

Failure modes: what breaks at 2 a.m.

Advanced means knowing how it fails. The two families fail differently:

Video: Your Stream Is Falling Behind: Backpressure, Lag, and Batching — SwetankLabs
What actually pages you at 2 a.m.: consumer lag, backpressure, skewed partitions, and buffer pileups — like a traffic jam on the ticker — plus how much spare capacity you need to drain it.

How to choose: a decision guide

There's no universally right answer, but the questions that pick one are refreshingly concrete:

And one connection worth making explicit: the durable log at the heart of both streaming architectures is the same Kafka you studied in the Message Queues lesson. There, it was a buffer between producers and consumers. Here, it's the system of record you compute on — and, in kappa's case, the thing you rewind when you need a second chance.

The cheat sheet, if you want it on one slide:

BatchStream
Data shapeBounded — the newspaper stackUnbounded — the live ticker
LatencyMinutes to hoursMilliseconds to seconds
ReprocessingNative superpower: re-run the jobReplay the retained log
Hard partsSkew, stragglers, schedulingEvent time, state growth, late data
Reach forSpark, MapReduceKafka + Flink

Video: Batch vs Streaming Explained for Data Engineering — TechLambda
Practical decision guide: when batch pipelines win, when streaming is worth the cost, and why hybrid architectures are common.

Takeaways

  1. Bounded vs unbounded is the core distinction: batch reads a finished dataset (the newspaper stack); streams process an endless flow (the live ticker).
  2. Batch gives you high throughput and complete, reprocessable answers at the cost of latency — Spark jobs on a schedule. Streams give you low latency at the cost of dealing with event time, windowing, and late data.
  3. Watermarks are the stream processor's answer to late events: a declared cut-off ("we've seen everything up to 9:00") that lets windows close — a heuristic, not a guarantee.
  4. Lambda merges batch and speed layers for freshness plus completeness, priced in duplicated logic and operational complexity. Kappa drops the batch layer entirely and replays a retained immutable log — simpler, but only as good as your retention.

Check your understanding

  1. In this lesson's analogy, the live news ticker represents...

    • The serving layer of a Lambda architecture
    • A bounded dataset, like yesterday's server logs
    • An unbounded data stream that never ends
    • A nightly Spark job that precomputes reports
  2. Some click events arrive four minutes after the window they belong to has closed. What do watermarks let the stream processor do?

    • Declare a cut-off time up to which events are assumed seen, so windows can close without waiting forever
    • Pause the event log until every straggler has arrived
    • Retroactively rewrite the results of closed windows for free
    • Guarantee that no event will ever arrive late
  3. A team running a Lambda architecture changes how a metric is computed. What must they do to stay correct?

    • Delete the serving layer and query the event log directly
    • Update and re-run both the batch job and the streaming job, so the two layers keep agreeing
    • Only update the batch layer — the serving layer backfills the speed layer
    • Only update the speed layer — the batch layer recomputes itself automatically
  4. Kappa architecture can fully replace a batch layer only when...

    • The stream processor runs on more machines than the batch cluster had
    • Events arrive in perfect timestamp order, so watermarks are unnecessary
    • The serving layer caches every historical query result
    • The immutable event log is retained long enough to replay the full history through new logic

Go deeper

Want to keep pulling this thread? These talks and tutorials go further than we did here:

Sources & further reading