1. What It Is
Stream processing turns unbounded event flows into real-time aggregates β windowed counts, joins, and alerts. We pick tumbling, sliding, or session windows based on what "current" means for the product.
What:
Real-time processing engines (Apache Flink, Spark Streaming) that aggregate and transform infinite message streams continuously.
Primary purpose:
Extracting analytical insights, metric counters, or fraud signals from event logs instantly as they occur.
Usually used for:
Real-time GPS coordinate ETAs, live payment fraud alerts, and sliding-window system health trackers.
2. Core Mental Model
Batch jobs answer "what happened yesterday?" Stream processors answer "what is happening now?" The hard part is not throughput β it is correctness under delay: events arrive out of order, clocks lie, and state must survive crashes. Problem #16 Distributed Stream Processing is the full Flink HLD; this concept covers the primitives both Flink and Kafka Streams share.
π°οΈ Event vs Processing Time
Event Time is when the update occurred on the client. Processing Time is when it reaches the server. Always process using Event Time.
π Watermark advancement
A temporal anchor signal: 'We assume all events with Event Time < T have arrived.' Advance watermarks to close windows and emit states.
ποΈ Stateful RocksDB Cache
Store running aggregate counters (e.g. click counts) inside local RocksDB key-value stores to enable micro-second local state updates.
In the room
Candidates confuse stream processing with batch. Say whether you need sub-second latency (Flink/Spark Streaming) or minute-level is fine. Watermarks and late-arriving events trip people up β mention allowed lateness if ordering matters.
3. Why It Matters in HLD
Stream processing turns unbounded event logs into real-time aggregates and alerts β Flink, Kafka Streams, Spark Streaming. Three lenses:
Needed When:
Designing real-time analytical systems where batch jobs (e.g. Spark/Hadoop running nightly) are too slow to drive immediate decisions.
Avoids:
Stale analytical dashboards, database locks from continuous aggregate writing, and lost telemetry insights.
Optimizes For:
Insight delivery speed, state consistency metrics, network resource conservation, and local storage utilization.
4. Architecture & Data Flow
Walk a streaming topology as interview steps. Step 1 β Source: consume from Kafka partition with offset tracking. Step 2 β Parse/transform: map, filter, enrich with join to reference data. Step 3 β Window: tumbling/sliding/session window for aggregates. Step 4 β Sink: write results to DB, cache, or downstream topic. Step 5 β Checkpoint: state backend snapshots for exactly-once recovery.
5. Key Characteristics
Event time vs processing time, window types, and stateful operators β we compare:
- Window types β how infinite streams are bounded for aggregation:
| Window Strategy | Operational Mechanic | Primary Use Case |
|---|---|---|
| Tumbling (Fixed-size) | Non-overlapping time brackets (e.g. every 5 minutes: [10:00, 10:05], [10:05, 10:10]). | Hourly active user counts, summary billing aggregations. |
| Sliding (Overlapping) | Fixed duration evaluating at regular steps (e.g. 5-min window rolling every 1 min). | Real-time system CPU spike trackers, sliding-window rate limiters. |
| Session (Dynamic) | Tiers grouped by user inactivity thresholds (e.g. closes after 30 mins idle). | Analyzing website visitor navigation paths or user gaming sessions. |
- Flink vs Kafka Streams β when to name which in an interview:
| Engine | Strength | Tradeoff |
|---|---|---|
| Apache Flink | Stateful stream processing with event-time windows, watermarks, and checkpoint-based exactly-once semantics across failures. |
|
| Kafka Streams |
|
|
Exactly-once processing (EOS) summary: At-least-once delivery + idempotent sinks is often enough. True EOS requires atomic checkpoint of operator state and output offsets together β Flink injects checkpoint barriers and replays from the last consistent snapshot on failure; Kafka Streams uses transactional writes with processing.guarantee=exactly_once_v2. Mention EOS when money or billing aggregates are at stake; at-least-once + idempotent counters suffices for dashboards.
In the room
Say event time vs processing time when discussing windows β interviewers probe whether you understand skewed clocks and late arrivals.
6. Strategic Tradeoffs
Real-time insight trades late-data handling and operational complexity β we state both:
| Benefit | Cost |
|---|---|
| Sub-second latency (evaluate each event as it arrives instead of waiting for batch jobs) | Stateful storage (open windows and counters live in local RocksDB or similar, adding memory and recovery complexity) |
| Correct ordering under delay (watermarks let late events land in the right window before it closes) | Memory for open windows (waiting for late data keeps window state in memory longer) |
7. Failure / Bottleneck Awareness
Late events, state explosion, and duplicate processing β we name mitigations:
Problem: A mobile client goes offline inside a subway. When it reconnects 1 hour later, it uploads coordinates. Because the event-time watermark advanced past that window, the streaming engine drops the late-arriving events, losing data.
Mitigation: Configure allowed lateness so closed windows accept late events for a grace period, or route very late events to a side output for manual reconciliation.
Problem: Session windows tracking user states without strict expiry limits accumulate infinite RocksDB states, eventually triggering JVM Out-Of-Memory crashes.
Mitigation: Set TTL on state entries, cap session idle timeouts, and checkpoint state to durable storage (e.g., S3) so workers can recover after restarts.
8. Common HLD Usage
Metrics dashboards, fraud detection, and live leaderboards use stream processing:
| Production System | Streaming Engine Selection | Architectural Rationale |
|---|---|---|
| Uber Ride ETA Calculator | Apache Flink (Streaming) | Continuous GPS coordinates telemetry streams must be aggregated inside tumbling windows to emit real-time driver ETA updates. |
| Financial Fraud Detector | Kafka Streams / Spark | Evaluating transactions within sliding 1-minute windows to spot rapid multi-charge card usage attempts. |
9. Decision Signals
Add stream processing when batch latency (hours) is too slow for the product SLA:
- You are designing telemetry monitors, ride ETA calculators, financial trade alerts, or sliding rate limiters.
- Your processing latency requirements are strictly under 5 seconds (invalidating overnight batch engines).
- You must evaluate event logs that arrive out-of-order due to variable mobile network latencies.
11. Deep Dive (Optional)
Stream Processing Fault Tolerance: Chandy-Lamport Distributed Checkpointing
To guarantee **Exactly-Once Processing** under node failures, Flink uses a variation of the **Chandy-Lamport algorithm** for asynchronous state checkpointing:
- The stream source periodically injects a **Checkpoint Barrier** (metadata packet) into the event stream.
- As the barrier flows through processing operators, each operator pauses its stream processing.
- The operator copies its active state (e.g. running click counters) asynchronously to a durable backup store (e.g. AWS S3).
- Once the operator finishes saving its checkpoint, it forwards the barrier down-stream.
- If a worker node crashes mid-stream, Flink rolls back *all* operators to the last successful checkpoint barrier index and replays logs from the Kafka event offsets.
This gives exactly-once processing semantics across failures: on crash, Flink rewinds to the last checkpoint and replays Kafka from the stored offsets.
Review
How helpful was this walkthrough?
Click a star to rate. We actively use this feedback to refine and update our system design content.
Discussion
Share your thoughts, ask questions, or help others.