0%
Data-Intensive Applications
Foundations of Data Systems
Distributed Data
Encoding and Evolution
Batch Processing
Data Quality and Governance
Operational Patterns
Lambda and Kappa Architectures
Why do data-intensive systems need two processing paths? Because no single path delivers both correctness and low latency. A batch job over the full dataset gives you accurate, complete results, but it takes hours to finish. A stream processor gives you answers in seconds, but those answers are approximate because streams see only recent events without full historical context.
This tension is not theoretical. Every analytics platform, every recommendation engine, and every fraud detection system faces it. Users want accurate historical reports and they want real-time dashboards. The same company that needs "total revenue this quarter" (accurate, can wait) also needs "revenue this minute for the live dashboard" (approximate, must be instant).
The Lambda Architecture, introduced by Nathan Marz in his book "Big Data: Principles and Best Practices of Scalable Real-Time Data Systems," solves this by running both paths simultaneously and merging their outputs.
The name "Lambda" comes from the Greek letter, used here to represent a function that produces the same output from either path.
The Three Layers
Batch layer: Every incoming event is appended to an immutable, append-only master dataset (think HDFS or S3). Periodically, a batch job (MapReduce, Spark) recomputes views over the entire dataset from scratch. The output is a set of batch views: precomputed query results that are complete and accurate. The downside is latency. A batch job over terabytes of data might take 2-6 hours, so the batch views are always stale by at least one batch cycle.
Speed layer: The same incoming events also flow into a stream processor (Storm, Flink, Spark Streaming). The speed layer computes real-time views that cover only the window since the last batch run completed. These views are approximate: the stream processor might count events that arrive out of order, double-count retries, or miss events during a restart. But the views are fresh, typically seconds behind real time. The speed layer's views are temporary: every time a batch run completes, the speed layer discards its state for the covered window and starts fresh.
Serving layer: Query-time logic merges the batch views and speed views to produce a complete answer. For a "total page views" query, the serving layer adds the batch count (accurate through 6 AM) to the speed count (approximate from 6 AM to now). The batch views provide correctness; the speed views fill the recency gap.
The serving layer must know the exact boundary (timestamp or offset) where batch coverage ends and speed coverage begins. Getting this boundary wrong means either a gap (missing events) or an overlap (double-counted events). In practice, the batch job writes a metadata record with the timestamp of the last event it processed. The serving layer reads this timestamp and queries the speed layer for events after that point.
The batch layer's key property is that it recomputes from scratch every cycle. If your batch job has a bug, you fix the code and re-run it over the full dataset. The immutable master dataset means you never lose raw data. This is the Lambda Architecture's strongest guarantee: eventual correctness through recomputation.
Why the Batch Layer Recomputes from Scratch
This sounds wasteful, but it is intentional. Incremental computation is fragile: if the logic changes, you cannot apply the new logic to already-processed data without reprocessing. Consider a simple example: your batch job counts unique visitors by country. You discover a bug where IP-to-country mapping was wrong for VPN users. With incremental computation, fixing this requires tracing through every previous result, identifying which counts were affected, and patching them individually. With scratch recomputation, you fix the IP-to-country mapping function, re-run the batch job, and the correct counts emerge automatically.
By recomputing from the immutable master dataset every cycle, the batch layer automatically incorporates code fixes, schema changes, and new derived fields. The cost is compute time and resources. The benefit is simplicity: you never need migration scripts or backfill jobs for the batch views.
This property becomes more valuable as your system grows. At 10 batch views, manual patching after a bug is tedious but feasible. At 500 batch views, recomputation is the only practical approach. It also simplifies onboarding: a new engineer does not need to understand the history of schema changes to reason about the current batch output. Every run starts clean.
The cost of recomputation is predictable and plannable. You know how long a batch cycle takes because you run it every few hours. If a Spark job over 10 TB takes 3 hours today, it will take approximately 3 hours next week (assuming similar data volume growth). This predictability lets operations teams plan maintenance windows, allocate cluster resources, and set expectations with stakeholders about when fixes will be visible in batch views.
A Concrete Example: Ad Click Counting
Consider a system that counts ad clicks for billing. Advertisers need accurate totals (they pay per click), and dashboards need real-time visibility (campaign managers want to see current spend).
The batch layer runs a Spark job every 4 hours over the full click log stored in S3. It performs several operations that require full-dataset access:
- Deduplicates clicks by click ID across the entire log (catching duplicates that arrive hours apart)
- Filters bot traffic using a trained ML classifier that evaluates click patterns across sessions
- Applies currency conversion at the day's official exchange rate
- Writes per-campaign totals to a serving database
These totals are authoritative for billing.
The speed layer runs a Flink job that reads clicks from Kafka in real time. It applies a simpler deduplication (in-memory set of recent click IDs, window of 10 minutes) and a basic bot filter (IP blocklist rather than a trained ML model, because the model takes too long to run per-event at real-time scale). The real-time counts are approximate: they might include some bot clicks that the batch classifier would catch, and they might miss deduplication for clicks that arrive more than 10 minutes apart. But the dashboard shows spend updating every second, which campaign managers need for budget pacing decisions like pausing a campaign that is burning through its daily budget too fast.
The serving layer responds to "how many clicks for campaign X today?" by returning the batch total (accurate through the last completed batch) plus the speed total (approximate for the window since the batch). Every 4 hours, the new batch run replaces the speed layer's estimates with exact counts, and the speed layer resets its window.
The timeline for a single day looks like this:
- 00:00: Batch run starts over yesterday's full click log
- 02:30: Batch run completes, writes exact counts for yesterday to serving DB
- 02:30-06:00: Speed layer covers only the gap since midnight
- 06:00: Next batch run starts (covering midnight to 06:00)
- 08:30: Batch run completes, speed layer resets to cover only 08:30 onward
- This cycle repeats every 4 hours
At any moment, the dashboard shows: batch total (exact, covering the last completed window) + speed total (approximate, covering events since then).
The Serving Layer Merge
The merge strategy depends on the query type. For additive metrics (counts, sums), the serving layer adds batch and speed results. For non-additive metrics (unique visitors, percentiles), the merge is more complex: you might use HyperLogLog sketches in the speed layer and exact counts in the batch layer, then take the batch result and add only the speed layer's estimate for the uncovered window.
The merge logic itself must be carefully designed. If the batch view covers events through timestamp T, the speed view must cover only events after T. If both cover the same window, you double-count. This coordination between layers adds operational complexity: the serving layer needs to know the exact boundary between batch and speed coverage, and that boundary changes every time a batch run completes.
There are two common merge strategies:
Time-boundary merge: The serving layer knows that the batch view covers events up to timestamp T. It queries the speed view for events after T. This works well for time-series data but requires accurate, synchronized timestamps across both layers.
Key-based merge: For key-value views (user profiles, product scores), the serving layer checks whether a key exists in the batch view. If yes, it uses the batch value. If no, it falls back to the speed view. This is simpler but assumes the batch view is always more authoritative than the speed view.
The choice of merge strategy affects how users perceive data quality. With time-boundary merge, users see a smooth transition: numbers gradually update in real time, then snap to exact values when the batch completes. With key-based merge, users see a more abrupt transition: a key might show approximate speed-layer values for hours, then suddenly switch to the exact batch value.
Written out, the time-boundary merge is small, and its smallness is deceptive:
The two functions differ in one way that decides the whole design. Clicks are additive, so the speed layer can hold a plain integer and the merge is +. Unique viewers are not, so the speed layer has to hold a structure that composes, and the batch layer has to hold one too even though it could have computed the exact answer. You cannot merge an exact count with a sketch: a viewer who appears in both windows is counted twice, and there is no correction because the batch layer threw away the identities. The mergeability of HyperLogLog, developed in architecture of real time analytics, is what makes the second function possible at all, and the price is that the batch layer gives up exactness to pay for it.
The other detail is batch.hwm. The boundary has to be a value the batch run publishes atomically with its output, not a wall-clock estimate of when the run covered. A batch job that reads "all events before midnight" but was delayed and actually ingested a few late arrivals stamped 00:00:03 will double-count those three seconds forever if the serving layer merges at a hard-coded midnight. Write the boundary into the batch view; never infer it.