Batch vs Streaming: Advanced
What you will be able to do
Lambda Architecture
Identify the three layers of Lambda architecture and explain why the constraints that motivated it have shifted.
The Three Layers
The Architecture in a Diagram
Why Marz Proposed It
| Constraint in 2011 | What It Forced | Lambda's Answer |
|---|---|---|
| Hadoop was the cheap correct engine | Batch was unavoidable for the canonical view | Use batch for the durable layer |
| Storm was the fast engine but tradeoffs were harsh | Streaming had at-most-once or at-least-once but rarely exactly-once | Use streaming only as an approximation; batch overwrites it |
| Storage was cheap; reprocessing was the cure for any bug | Recomputing from the master log was always available | The batch layer is a complete recompute on every run |
| Serving had to be fast and merge two sources | A serving layer that combined views was needed | HBase or Cassandra plus a cache, with merge logic |
Why Lambda Became Unfashionable
- Correctness is guaranteed by the batch layer overwriting the speed layer
- Streaming can be approximate, opening up cheaper engines
- The master event log is durable; bugs are fixed by reprocessing
- Each layer has clear semantics and a clear lifecycle
- Two implementations of the same logic in two languages
- Operational complexity doubles: two pipelines, two failure modes
- Drift between layers is real and hard to explain to consumers
- Serving layer merge logic adds latency and edge cases at the boundary
What Lambda Got Right
- ▸The master event log is the source of truth; everything else is a derived view
- ▸Recomputability is a first-class property; bugs are fixed by reprocessing, not patching
- ▸Serving is a separate concern from processing
- ▸Different freshness tiers can coexist on the same data, served from different views
- Keep the immutable event log as the source of truth, regardless of architecture
- Preserve the recompute-from-log property when migrating away from Lambda
- Document why Lambda was chosen historically; the constraints may inform the next architecture
- Build a new Lambda system in 2026 unless the streaming layer must be approximate
- Discard the master log when migrating away from Lambda; it is the most durable piece
- Treat Lambda as obsolete; its principles are still the right framing
Two ways data moves: batch processes a whole chunk on a schedule (accurate, delayed); streaming processes each event as it arrives (fast, continuous). The latency SLA decides which.
Kappa: Stream Only, Batch Replay
Apply Kappa architecture: stream-only with batch as replay, and explain what is given up in storage retention to gain a single codebase.
The Architecture in a Diagram
What Kappa Gives Up
| Lambda Property | Kappa's Tradeoff |
|---|---|
| Storage retention can be short; only the most recent window is in the log | Kappa requires keeping the full event log indefinitely or paying for cold replay |
| The batch layer can use cheap correct engines | Kappa pays streaming compute prices for everything, even the daily aggregate |
| Reprocessing happens automatically every batch run | Kappa requires explicit replay jobs and parallel materialized views |
| Approximations in the speed layer are acceptable | Kappa demands correctness in streaming; exactly-once is not optional |
Why Kappa Became the Default
Replay as the Universal Backfill
What Kappa Cannot Do
- Source data is naturally an event log (Kafka, Pulsar, CDC stream)
- Streaming engine supports exactly-once for the sinks in use
- Log retention is affordable; tiered storage is in place
- Single codebase is operationally cheaper than two codebases
- Source data is a periodic snapshot (no log to replay)
- Workload is rare and reading the log every replay is wasteful
- Streaming exactly-once is hard or unsupported for the sink
- Cost-per-row is the dominant constraint, not freshness
The Storage Tradeoff
- ▸Source is a Kafka topic, Pulsar topic, or CDC stream from the start
- ▸Exactly-once semantics are supported by the engine and the sink
- ▸Tiered storage or another long-retention mechanism is affordable
- ▸The team has streaming engineering capacity; one pipeline is one operational profile
Unified Engines: Where Lines Blur
Identify which aspects of batch and streaming are unified by modern engines and which still differ at runtime.
What the Unified Engines Unify
| Layer | Spark | Flink | Beam |
|---|---|---|---|
| API | Same DataFrame for batch and streaming | DataStream API; bounded stream is batch | Single Pipeline; runtime selects the engine |
| Trigger model | processingTime, once, continuous | Time-based, count-based, custom | Trigger objects abstracted across runners |
| Watermarks | withWatermark on event-time columns | Native watermark assigners | WatermarkStrategy at the source |
| State | Checkpointed RocksDB or HDFS | RocksDB local plus async snapshots | Runner-specific state backend |
What the Unified Engines Still Distinguish
- Whether the engine reads bounded or unbounded input
- Whether the trigger fires once or repeatedly
- Whether state is in-memory for one run or persistent across runs
- Whether the engine runs to completion or runs forever
- Cost: continuous compute vs. on-demand compute
- Failure mode: rerun partition vs. checkpoint-and-replay
- Observability: binary success vs. lag and latency percentiles
- Deployment: redeploy at next run vs. drain-and-replace mid-flight
Spark Structured Streaming as a Concrete Case
Where the Line Still Matters
| Concern | Why the Line Still Matters |
|---|---|
| Cost budgeting | Streaming and batch have order-of-magnitude different costs even at identical logic |
| On-call rotation | Streaming requires lag-based alerting; batch requires schedule-based alerting |
| Schema migration | Streaming requires drain-and-replace; batch swaps at the next run boundary |
| Backfill mechanics | Streaming replays from offsets; batch reruns by date partition |
| Failure-recovery testing | Streaming needs chaos tests during runs; batch needs partition replay tests |
How Senior Engineers Use Unified Engines
- ▸Can this batch pipeline run as micro-batch without a code change? Yes, with a config change.
- ▸If a team made this streaming pipeline batch for cost reasons, what would the team lose? The unified API lets both modes be run and compared.
- ▸How does the same logic behave under streaming versus batch failure modes? Run both and compare.
- ▸Can a streaming pipeline be backfilled by running its code as a one-shot batch? Yes, the unified API supports it.
- Write logic against a unified engine API so rhythm changes are config changes, not rewrites
- Test both batch and streaming runtimes for any pipeline that might graduate between them
- Keep observability separate per rhythm; unified API does not unify lag and schedule semantics
- Treat the unified API as a guarantee of identical operational behavior
- Pick a streaming-only engine for a workload that may always be batch
- Skip the cost conversation because the API is the same; runtime cost differs by an order of magnitude
Per-Node Freshness Tier Analysis
Annotate every node in a pipeline with an explicit freshness tier and identify mismatches that produce hidden cost or unmet consumer expectations.
Tiers Per Layer in a Layered Pipeline
| Layer | Typical Freshness Tier | Why |
|---|---|---|
| Source | Continuous (what the producer emits) | Cannot be tightened by the pipeline; bound is set upstream |
| Raw landing | Tier 2 to 3 (under 15 min to under 2 hr) | Lags source by ingestion overhead; cheap to keep tight |
| Curated | Tier 3 to 4 (under 2 hr to daily) | Joins and aggregations are expensive; refreshed when consumers actually read |
| Serving | Tier 1 to 4 (per consumer) | Consumer-facing; tier matches the specific consumer's need |
Why Mixing Tiers Is Fine
The Tier Label as a Diagram Element
How to Pick a Tier per Node
- ▸Start at the consumer. What freshness does each consumer actually need?
- ▸Walk backward. Each upstream node must be at least as fresh as its strictest downstream consumer.
- ▸Allow upstream nodes to be tighter if they have other consumers with stricter needs.
- ▸Allow upstream nodes to be looser if no downstream consumer reads them at the looser cadence.
- ▸Document each node's tier; mismatches between adjacent tiers are the most common bug.
Tier Mismatches as a Failure Mode
| Tier Mismatch | Symptom | Fix |
|---|---|---|
| Upstream daily, downstream hourly | Hourly view does not change between daily upstream runs | Tighten upstream cadence or relax downstream tier |
| Upstream tier-2 streaming, downstream nightly batch | Streaming work is wasted; nightly only sees what daily batch would see | Either consumer reads streaming directly, or upstream is downgraded to nightly |
| Two consumers at different tiers reading the same node | One consumer overpays for freshness, or the other underreceives | Split into two serving nodes at the appropriate tiers |
| Source faster than the rest of the pipeline | End-to-end latency floor is much higher than source latency | Tighten the slowest node; the latency is the max of all nodes |
The Cost Story Per Tier
- Every node refreshes at the strictest consumer's tier
- Cost is dominated by the strictest tier multiplied by every node
- Consumers with looser needs overpay for unused freshness
- Architecture is simpler but more expensive
- Each node refreshes at the tier its downstream consumers actually need
- Cost is the sum across nodes, each at its appropriate tier
- Consumers pay for the freshness they read; no excess
- Architecture is more complex but operates at minimum cost
Lambda to Kappa Worked Example
Walk through a Lambda-to-Kappa migration on a real workload and name what changes in code, storage, and operations.
The Lambda Starting Point
The Kappa Redesign
What Changes in Code
| Concern | Lambda | Kappa |
|---|---|---|
| Aggregation logic | Implemented twice: Spark Scala and Storm Java | Implemented once in Flink Java |
| Window semantics | Daily windows in batch; rolling windows in speed | Single windowing model in Flink with watermarks |
| Idempotency | Batch overwrites the day; speed approximates | Flink exactly-once with Iceberg ACID transactions |
| Backfill | Re-run Spark for the date range | Replay Flink from a chosen offset, write to parallel table |
| Schema migration | Coordinate two codebases; redeploy both | One Flink redeploy with state migration |
What Changes in Storage
What Changes in Operations
| Operational Concern | Lambda | Kappa |
|---|---|---|
| On-call rotation | Two rotations: batch on-call, streaming on-call | One rotation: streaming pipeline plus replay jobs |
| Failure recovery | Batch reruns the night; streaming restarts from checkpoint | Streaming restarts from checkpoint; backfill is a replay |
| Drift investigation | Reconcile batch vs speed; trace the discrepancy | No drift exists; eliminated by single source of truth |
| Cost attribution | Two cost centers (batch cluster, speed cluster) | One cost center; cost per pipeline is direct |
| Deployment cadence | Two pipelines, two release cycles | One pipeline, one release cycle |
What Stays the Same
The Migration Path
- ▸Confirm exactly-once semantics in the streaming engine for the existing sinks
- ▸Extend Kafka retention to cover the longest backfill window the team needs
- ▸Build the Kappa pipeline alongside Lambda, writing to a parallel materialized view
- ▸Validate the Kappa view matches the Lambda merged view within an acceptable tolerance
- ▸Migrate consumers one by one; the Lambda system runs alongside until the last consumer has cut over
- ▸Retire the Lambda batch layer and speed layer; the Kappa pipeline owns the workload
When the Migration Fails
- Two codebases (Spark + Storm)
- Drift between batch and speed layers
- Two on-call rotations
- Backfill is rerun-the-night
- Storage in HDFS plus HBase
- 0.4 percent unexplained drift
- One codebase (Flink)
- Single source of truth; no drift
- One on-call rotation
- Backfill is replay from offset
- Storage in Kafka tiered plus Iceberg
- Drift eliminated by single layer
> A Series E retail platform inherited a Lambda content engagement pipeline from its founding-team era. Two codebases, two on-calls, 0.4 percent drift between layers, and a CFO who asked last week why the data infrastructure bill grew 60 percent year over year. The new principal data engineer is asked to redesign the system, write a migration plan, and make the cost shape defensible.
Lambda, Kappa, and unified engines: architectures live or die on freshness tier discipline
- Category
- Pipeline Architecture
- Difficulty
- advanced
- Duration
- 35 minutes
- Challenges
- 0 hands-on challenges
Topics covered: Lambda Architecture, Kappa: Stream Only, Batch Replay, Unified Engines: Where Lines Blur, Per-Node Freshness Tier Analysis, Lambda to Kappa Worked Example
Lesson Sections
- Lambda Architecture (concepts: paLambdaArch)
Lambda architecture is the first widely adopted attempt to combine batch and streaming in one system. Nathan Marz proposed it around 2011 in his book Big Data, drawing on his experience at Twitter and BackType. The motivation was specific to the era: batch frameworks (Hadoop MapReduce) were correct but slow; stream frameworks (Storm) were fast but produced approximate results. Lambda combined the two, using batch for the durable correct view and streaming for the live approximate view. Both laye
- Kappa: Stream Only, Batch Replay (concepts: paKappaArch)
Kappa architecture, proposed by Jay Kreps in 2014, is the answer to Lambda's two-codebase problem. The idea is simple: keep only the streaming layer. The event log is the source of truth, the streaming pipeline produces the canonical view, and batch becomes a special case (replaying the event log through the same streaming pipeline) rather than a separate codebase. One implementation of the logic, one operational profile, one set of failure modes. The simplification is real, and Kappa has become
- Unified Engines: Where Lines Blur (concepts: paBatchVsStreaming)
The cleanest version of Kappa requires an engine that runs the same code in batch and streaming modes. Modern engines have moved toward this ideal. Spark Structured Streaming exposes a unified DataFrame API where the same query can run as a batch job, a micro-batch streaming job, or a continuous streaming job by changing one configuration. Apache Flink runs streaming as the default and batch as a special case (a bounded stream). Apache Beam abstracts both into a single programming model. The con
- Per-Node Freshness Tier Analysis (concepts: paBatchVsStreaming)
A single pipeline rarely needs one freshness tier across every node. The source might produce events continuously. The raw landing layer might lag the source by seconds. The curated layer might rebuild hourly. The serving layer might refresh on a per-consumer schedule. Treating the entire pipeline as one tier (the strictest one) overbuilds most nodes; treating it as the loosest tier underbuilds the consumer-facing edge. Senior engineers tier each node explicitly and label it on the architecture
- Lambda to Kappa Worked Example (concepts: paKappaArch)
The synthesis exercise walks through a real-shaped migration: a workload originally designed as Lambda, redesigned as Kappa, with explicit notes on what changes in code, in storage, and in operations. The example is a streaming media company's content engagement pipeline. The exercise shows that the migration is not a rewrite; it is a careful retirement of the batch layer and a tightening of the streaming layer, with the immutable event log surviving as the architectural anchor. The Lambda Start