Lambda Architecture
also called Batch and Speed Layer, Dual-Path Analytics
Running the same computation as both a scheduled recomputation over all retained history and a streaming job over recent events, so that any defect in the fast path is erased by the next full recomputation.
A stream processor's output is only ever as good as the code that was running when the event went past. Ship a bug on Tuesday, find it on Friday, and the aggregate for those three days is wrong in a store that has no memory of the inputs that produced it. There is no edit path, because the derivation was applied once and the inputs were forgotten.
Lambda, as Nathan Marz described it from 2011 onwards, answers that with redundancy rather than cleverness. A batch layer recomputes every view from the immutable master dataset on a schedule. A speed layer computes the same views incrementally over the events that have arrived since the last batch run. A serving layer answers queries by merging the two.
The structure looks wasteful until you notice what the schedule buys: the lifetime of a speed-layer defect is bounded by one batch cycle. Whatever the fast path got wrong today is discarded tomorrow, because tomorrow's answer is recomputed from the raw events rather than carried forward.
Why it matters
Most arguments about Lambda are about freshness, and freshness is not what it is for. It is a correction mechanism with a bounded worst case, and it is the only architecture in this family that assumes the engineers writing the streaming logic will be wrong.
That assumption is worth money in two places. Where a number is signed by a person, such as a monthly revenue figure or an advertising settlement, the signer needs a figure that is reproducible from source data and does not change after publication. Where the streaming logic is changed frequently, the batch recomputation is the regression test that runs against production data.
Implementation patterns
- An immutable master dataset, appended only, partitioned by event time, on object storage. This is the real asset; both layers are derived and disposable.
- An event-time cut published as data. The batch job records the watermark it covered, and the serving layer reads that value to decide where to switch views. A processing-time cut, such as "the batch ran at 02:00", double-counts or drops every late event near the boundary.
- Identically keyed views, so the merge is a lookup in both and a combine, not a join.
- An idempotent merge, either because the metric is a max or a set union, or because records carry event ids and the serving layer deduplicates.
- One codebase, two runners. The transformation is a library with no I/O, invoked by a streaming runner over the log and a batch runner over object storage. This removes the drift that makes classical Lambda expensive while keeping both execution paths.
- A speed-layer state TTL slightly longer than the batch interval, so the fast path always covers the gap even when a batch run is late.
Industry example
Marz's original formulation came out of building analytics on an immutable event log and was written up in his 2011 post on beating the CAP theorem and later in book form. The shape recurs wherever one event stream has two audiences with different requirements.
A live video platform at Hotstar or JioCinema scale has exactly that split: a concurrency counter must be on an operations wall inside a few seconds, and the advertising delivery number for the same match has to be exact, reproducible and agreed with a counterparty the next morning. The first requirement tolerates a figure that is 2% out. The second does not tolerate being 0.1% out. Running one path and asking it to serve both means the strict consumer inherits the loose consumer's correction model.
Failure scenarios
- The processing-time cut. Events near the boundary appear in both views or in neither, typically a fraction of a percent, which survives spot checks and surfaces at month end.
- The batch job quietly stops. The speed layer keeps answering, so dashboards look alive while the correction mechanism is gone. Staleness and drift accumulate until the speed layer's state expires, and then the number collapses without warning.
- Logic drift between the two implementations. Discovered during an argument about which number is right, usually in a meeting rather than by a test.
- The merge becomes the complex part. It is the least tested code in the system and the only code that has to understand both views' edge cases.
- Batch interval longer than speed-layer retention, which leaves an uncovered window nobody notices until a batch run is delayed.
Trade-offs
| Choose | Gains | Pays |
|---|---|---|
| Two paths, two codebases | A bounded correction window and a signable number | Duplicated logic, drift, a merge layer, double the change cost per metric |
| Two runners, one library | Same correction window without the drift | A transformation layer with no I/O, which constrains how you write the logic |
| One streaming path with replay | One implementation | Hot log retention for the full correction window, keyed sinks, effect isolation |
When not to use it
If the log can be replayed across the whole correction window and the sinks are keyed and the side effects are separated from the computation, delete the batch layer. Those three conditions give a single path the same repair capability, and nothing else in Lambda is worth the second codebase.
Also skip it when nobody consumes an exact number. A product analytics dashboard with no auditor and no settlement has no buyer for the batch layer, and a team of four engineers maintaining two paths for one metric will maintain neither well. The honest default for a small team is one path plus a documented replay drill that has actually been run.
Interview question
Q: Your company computes daily revenue twice, once in a nightly batch job and once in a streaming job, and the two disagree by 0.3%. Which one do you trust, and what would you change?
What a strong answer covers: trust neither until you know where the cut is, because the most likely cause of a small systematic gap is a boundary defined in processing time on one side and event time on the other. Then: check the speed layer's late-data policy against the batch job's input filter, reconcile on a single day by event id rather than on totals, and make the batch job publish the watermark it covered so the merge stops guessing. A strong answer also says what it would do structurally, which is collapse the two implementations onto one transformation library before trying to explain a 0.3% gap between two separately written ones.
Quick check
Quiz: In Lambda, what bounds the damage a bug in the streaming path can do? The batch interval: the next full recomputation from the master dataset discards whatever the speed layer got wrong.
Flashcard: Why must the batch and speed layers be divided by an event-time cut rather than a wall-clock cut? Because late events whose event time falls before the cut but which arrive after the batch ran land in both views or in neither, which shows up as a sub-percent systematic error.