Core Concept

Stream Processing Basics

Stream processing computes aggregations and alerts over continuous event flows, using event-time windows and watermarks to handle out-of-order arrivals.


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.

Loading...

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 StrategyOperational MechanicPrimary 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:
EngineStrengthTradeoff
Apache FlinkStateful stream processing with event-time windows, watermarks, and checkpoint-based exactly-once semantics across failures.
  • Separate cluster to operate
  • best for complex joins, CEP, and sub-second analytics at scale.
Kafka Streams
  • Library embedded in your JVM service β€” no extra cluster. Reads/writes Kafka topics directly
  • EOS via idempotent producer + transactional commits.
  • Tied to Kafka as source/sink
  • less suited to multi-source joins or heavy stateful CEP than Flink.

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:

BenefitCost
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:

🌩️ The Late-Event Window Closed Starvation

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.

🐒 State Size Memory Bloat (OOM crashes)

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 SystemStreaming Engine SelectionArchitectural Rationale
Uber Ride ETA CalculatorApache Flink (Streaming)Continuous GPS coordinates telemetry streams must be aggregated inside tumbling windows to emit real-time driver ETA updates.
Financial Fraud DetectorKafka Streams / SparkEvaluating 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:

🎯 Think Stream Processing When:
  • 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:

  1. The stream source periodically injects a **Checkpoint Barrier** (metadata packet) into the event stream.
  2. As the barrier flows through processing operators, each operator pauses its stream processing.
  3. The operator copies its active state (e.g. running click counters) asynchronously to a durable backup store (e.g. AWS S3).
  4. Once the operator finishes saving its checkpoint, it forwards the barrier down-stream.
  5. 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

Help Us Improve

How helpful was this walkthrough?

Click a star to rate. We actively use this feedback to refine and update our system design content.

Placeholder
Optional but highly appreciated!

Discussion

Share your thoughts, ask questions, or help others.

Loading comments...