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.
8 minConcept guideReference-informed · independently authored
01

Lesson spine

What you need to understand.

Stream processing computes continuously over unbounded, imperfectly ordered events using explicit time and recovery semantics.

01

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.

02

Windows

Tumbling windows do not overlap, sliding windows overlap by a step, and session windows close after inactivity.

03

Event time and watermarks

Use event timestamps for business truth and a watermark to estimate when a window is complete despite late arrivals.

04

Stateful processing

Partition by key, update durable local state, and emit versioned results to an idempotent sink.

05

Checkpoints

Align input offsets with operator state snapshots so recovery replays only work after the last consistent checkpoint.

06

Late data

Drop by policy, update an open result, or emit a correction; the product must define how old truth may change.

02

Before the boxes

Frame the decision.

Outcome

What must work

Design windows, watermarks, state, checkpoints, reprocessing, and late-event policy.

Scale

What changes the design

Events/s · out-of-order delay · state size · window · checkpoint · recovery

Boundary

What owns the truth

Identify the component that commits authoritative state, then separate synchronous confirmation from derived work.

Non-goal

What stays simple

Do not add global coordination, multi-region writes, or a specialized store until a requirement earns the complexity.

03

Architecture map

Trace ownership, not just traffic.

Stream processing and real-time aggregation · concept mechanism

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.

  1. 01

    Sources — Emit timestamped events Define the output contract before moving to the next owner.

  2. 02

    Log — Preserves partition order Define the output contract before moving to the next owner.

  3. 03

    Processor — Updates keyed state Define the output contract before moving to the next owner.

  4. 04

    Watermark — Estimates completeness Define the output contract before moving to the next owner.

  5. 05

    Window — Emits aggregates Define the output contract before moving to the next owner.

  6. 06

    Checkpoint — Captures state Define the output contract before moving to the next owner.

  7. 07

    Sink — Publishes idempotent results Confirm the result and emit the evidence needed to reconcile it.

04

Decision table

Make the trade-offs explicit.

DecisionDefensible positionCost to acknowledge
Primary mechanismEvent 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. asyncKeep 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. scaledBegin 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.
05

Failure review

Design the recovery path.

DetectBoundRetry safelyReconcileLearn

Topic-specific risk

Late events and inconsistent checkpoints create wrong or duplicate output.

Response

Persist 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.

Response

Use 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.

Response

Expose queue depth and hot-key share, apply backpressure, isolate tenants, and degrade optional work before correctness.

06

Evidence + level bar

Prove the design can be operated.

Core signals

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.

Mid-level

Complete and clear

Finish the happy path, identify the state owner, choose reasonable building blocks, and explain one scale mechanism.

Senior

Trade-offs and failure

Separate read and write paths, define consistency, explain partitioning, and make duplicate or partial failure safe.

Staff+

Evolution and operations

Discuss multi-region boundaries, migration, tenant isolation, capacity, observability, and how the architecture changes over time.

07

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?

  1. For Stream processing and real-time aggregation, where is the correctness boundary and which failure would you test first?
  2. Which component owns committed truth, and what event or response proves the commit?
  3. Where is the first scaling or coordination bottleneck under the stated envelope?
  4. What happens after an ambiguous timeout or duplicate operation?
  5. Which complexity would you remove at one hundredth of the scale?