05Design patternsStream processing and real-time aggregation
05 · Reusable pattern
Stream processing and real-time aggregation
Design windows, watermarks, state, checkpoints, reprocessing, and late-event policy.Lesson spine
What you need to understand.
Stream processing computes continuously over unbounded, imperfectly ordered events using explicit time and recovery semantics.
Batch vs stream
Batch waits for a bounded collection; streams update state as events arrive. Many systems use both for different latency and correction needs.
Windows
Tumbling windows do not overlap, sliding windows overlap by a step, and session windows close after inactivity.
Event time and watermarks
Use event timestamps for business truth and a watermark to estimate when a window is complete despite late arrivals.
Stateful processing
Partition by key, update durable local state, and emit versioned results to an idempotent sink.
Checkpoints
Align input offsets with operator state snapshots so recovery replays only work after the last consistent checkpoint.
Late data
Drop by policy, update an open result, or emit a correction; the product must define how old truth may change.
Before the boxes
Frame the decision.
What must work
Design windows, watermarks, state, checkpoints, reprocessing, and late-event policy.
What changes the design
Events/s · out-of-order delay · state size · window · checkpoint · recovery
What owns the truth
Identify the component that commits authoritative state, then separate synchronous confirmation from derived work.
What stays simple
Do not add global coordination, multi-region writes, or a specialized store until a requirement earns the complexity.
Architecture map
Trace ownership, not just traffic.
Follow the decision from left to right. Every arrow should have a reason.
Emit timestamped events
Preserves partition order
Updates keyed state
Estimates completeness
Emits aggregates
Captures state
Publishes idempotent results
Walk one representative request across every arrow. Say whether the handoff is synchronous or asynchronous, what identity makes a retry safe, and which step changes authoritative state.
- 01
Sources — Emit timestamped events Define the output contract before moving to the next owner.
- 02
Log — Preserves partition order Define the output contract before moving to the next owner.
- 03
Processor — Updates keyed state Define the output contract before moving to the next owner.
- 04
Watermark — Estimates completeness Define the output contract before moving to the next owner.
- 05
Window — Emits aggregates Define the output contract before moving to the next owner.
- 06
Checkpoint — Captures state Define the output contract before moving to the next owner.
- 07
Sink — Publishes idempotent results Confirm the result and emit the evidence needed to reconcile it.
Decision table
Make the trade-offs explicit.
| Decision | Defensible position | Cost to acknowledge |
|---|---|---|
| Primary mechanism | Event time reflects reality but requires watermarks and correction; processing time is simpler. | The stronger guarantee usually adds coordination, latency, state, or operational work. |
| Sync vs. async | Keep only correctness-critical confirmation synchronous. Move derived views, notifications, analytics, and cleanup behind a durable boundary. | Async work needs idempotency, lag monitoring, replay, and a product definition for partial completion. |
| Simple vs. scaled | Begin with one logical owner and a clear API. Partition or replicate only the resource proven to be the first bottleneck. | Migration requires stable identities, versioned contracts, backfill, and a rollback path. |
Failure review
Design the recovery path.
Topic-specific risk
Late events and inconsistent checkpoints create wrong or duplicate output.
ResponsePersist enough identity and state to distinguish retry, resume, compensation, and operator repair.
Dependency timeout
A timeout is ambiguous: the remote side may have failed, succeeded, or still be running.
ResponseUse deadlines, bounded backoff with jitter, idempotency keys, and a status or reconciliation path.
Overload or skew
Average capacity can look healthy while a tenant, key, partition, region, or expensive request saturates one owner.
ResponseExpose queue depth and hot-key share, apply backpressure, isolate tenants, and degrade optional work before correctness.
Evidence + level bar
Prove the design can be operated.
Health of the promise
Measure user-visible latency or freshness, correctness drift, saturation, retry volume, and time to recover. Alert on the failed promise—not only CPU.
Complete and clear
Finish the happy path, identify the state owner, choose reasonable building blocks, and explain one scale mechanism.
Trade-offs and failure
Separate read and write paths, define consistency, explain partitioning, and make duplicate or partial failure safe.
Evolution and operations
Discuss multi-region boundaries, migration, tenant isolation, capacity, observability, and how the architecture changes over time.
Interview language
Open the deep dive with a claim.
“For Stream processing and real-time aggregation, the decision I want to make explicit is this: Event time reflects reality but requires watermarks and correction; processing time is simpler. I’ll trace the state-changing path first, show where the result becomes durable, then test the design against the highest-risk failure and our target scale.”
08 · Retrieval check
Can you defend it without the page?
- For Stream processing and real-time aggregation, where is the correctness boundary and which failure would you test first?
- Which component owns committed truth, and what event or response proves the commit?
- Where is the first scaling or coordination bottleneck under the stated envelope?
- What happens after an ambiguous timeout or duplicate operation?
- Which complexity would you remove at one hundredth of the scale?