🌊 Category 02 • Real-Time Event Pipelines

Streaming: Kafka, Spark Streaming & Events

Mastering low-latency, event-driven data architectures. Learn Apache Kafka internals, producer/consumer group protocols, schema evolution, Spark Structured Streaming micro-batches, stateful RocksDB checkpointing, and exactly-once delivery.

1. Real-Time Streaming Overview: Batch vs Streaming

Paradigm Shift

Traditional batch architectures process historical data in discrete windows (hourly, daily), leading to high latency and high peak cluster loads. Modern event-driven platforms shift from polling databases to continuous stream processing where data is consumed as it happens. We contrast the classical Lambda Architecture (separate batch and speed layers) with the unified Kappa Architecture (single append-only log with stream processors replaying historical offsets).

Batch Processing (Bounded)

Operates on fixed, finite datasets. High latency (hours), heavy disk I/O, optimal for large historical reconciliations and ML training.

Stream Processing (Unbounded)

Processes continuous, infinite event feeds with sub-second latencies using tumbling, sliding, or session watermarked windows.

Topic 2.1

Apache Kafka Fundamentals & Cluster Architecture

Distributed Systems

Apache Kafka is a distributed, partitioned, replicated commit log service. It provides high throughput, fault tolerance, and horizontal scaling by sharding topics across multiple broker instances.

Brokers, Clusters & KRaft

Kafka brokers store append-only segment files. Modern Kafka 3.x+ uses KRaft (Kafka Raft consensus) instead of external Zookeeper.

Topics, Partitions & Offsets

Topics are logical streams divided into partitions. Each record within a partition receives a sequential, immutable integer offset.

Replication & In-Sync Replicas (ISR)

Partitions have 1 Leader and N Followers. The In-Sync Replicas (ISR) set includes all followers caught up with the leader offset.

Kafka Connect & CDC

Declarative framework connecting databases to Kafka without code. Source connectors (Debezium for MySQL/PostgreSQL CDC) and Sink connectors (S3, Snowflake).

Topic 2.2

Producers, Consumer Groups & Rebalance Protocol

Core Mechanics

Producers route messages using key hashing, while Consumer Groups allow multiple workers to read from partitions in parallel. Managing rebalance storms and offset commits is crucial for stability.

Producer ACKs & Idempotence

Setting acks=all (or -1) combined with enable.idempotence=true ensures zero duplicate messages and zero data loss on broker retries.

Consumer Groups & Partitions

Each partition is read by exactly one consumer within a group. If you have 12 partitions, you can run at most 12 concurrent active consumers.

Consumer Rebalance Protocol

Cooperative Sticky Assignor prevents "stop-the-world" pauses when new consumer pods spin up in Kubernetes, reassigning only migrating partitions.

Manual vs Auto Offset Commits

Disable enable.auto.commit to avoid data loss. Commit offsets synchronously (commitSync) only after downstream storage succeeds.

📄 resilient_kafka_producer.py
import json
from confluent_kafka import Producer

def delivery_callback(err, msg):
    """Callback triggered on broker ACK or failure."""
    if err is not None:
        print(f"Message delivery failed: {err}")
    else:
        print(f"Message delivered to {msg.topic()} [{msg.partition()}] at offset {msg.offset()}")

conf = {
    'bootstrap.servers': 'msk-cluster.aws.internal:9092',
    'client.id': 'market-data-producer',
    'acks': 'all',                          # Strong durability
    'enable.idempotence': True,             # Prevents duplicates on retry
    'compression.type': 'snappy',          # Low latency compression
    'linger.ms': 20,                        # Micro-batching for throughput
    'batch.num.messages': 5000,
}

producer = Producer(conf)

def publish_market_tick(ticker: str, price: float, volume: int):
    payload = {"ticker": ticker, "price": price, "volume": volume}
    # Key hashing on ticker ensures all ticks for 'AAPL' land on the exact same partition
    producer.produce(
        topic='market-ticks',
        key=ticker.encode('utf-8'),
        value=json.dumps(payload).encode('utf-8'),
        callback=delivery_callback
    )
    producer.poll(0)

# Flush remaining events on shutdown
producer.flush()
Topic 2.3

Schema Management: Schema Registry & Avro/Protobuf

Data Contracts

Without schema enforcement, upstream producer changes break downstream consumer pipelines. Confluent Schema Registry / AWS Glue Schema Registry enforces explicit schemas for Avro, Protobuf, or JSON Schema.

BACKWARD Compatibility

Consumers using the new schema can read records produced with the previous schema. New fields must have default values.

FORWARD Compatibility

Consumers using the old schema can read records produced with the new schema. Deleted fields must have had default values.

FULL Compatibility

Both backward and forward compatible. Ensures consumers can be upgraded before or after producers with zero pipeline downtime.

Wire Format Overhead

Avro embeds a 5-byte Magic Byte + Schema ID in message headers rather than full JSON keys, slashing network payload size by 65%.

Topic 2.4

Spark Structured Streaming & Watermarking

Real-Time Aggregations

Spark Structured Streaming unifies batch and streaming APIs under Catalyst. It processes streams as an unbounded table, maintaining state in memory or RocksDB and handling late-arriving data via watermarks.

Event-Time vs Processing-Time

Aggregations group on event_time embedded in the message payload rather than when the Spark executor received it.

Watermarking

Tells Spark how late an event can arrive before being discarded (e.g. withWatermark("event_time", "10 minutes")). Prunes state store memory.

Output Modes

Append: Emits rows only after watermarked window closes. Update: Emits rows every time aggregate changes. Complete: Emits entire table.

Stateful Processing & RocksDB

Configure RocksDB state store provider (spark.sql.streaming.stateStore.providerClass) to avoid JVM garbage collection pauses on large stream states.

📄 spark_watermarked_streaming.py
from pyspark.sql import SparkSession
from pyspark.sql import functions as F
from pyspark.sql.types import StructType, StringType, DoubleType, LongType, TimestampType

spark = SparkSession.builder \
    .appName("KafkaMarketDataStreaming") \
    .config("spark.sql.streaming.stateStore.providerClass", 
            "org.apache.spark.sql.execution.streaming.state.RocksDBStateStoreProvider") \
    .getOrCreate()

# Schema definition
tick_schema = StructType() \
    .add("ticker", StringType()) \
    .add("price", DoubleType()) \
    .add("volume", LongType()) \
    .add("trade_time", TimestampType())

# Read streaming feed from Kafka
kafka_stream = spark.readStream \
    .format("kafka") \
    .option("kafka.bootstrap.servers", "msk-cluster.aws.internal:9092") \
    .option("subscribe", "market-ticks") \
    .option("startingOffsets", "latest") \
    .load()

# Parse JSON payload
ticks_df = kafka_stream.select(
    F.from_json(F.col("value").cast("string"), tick_schema).alias("data")
).select("data.*")

# 5-minute tumbling candlestick aggregation with 10-minute watermark
ohlc_candles = ticks_df \
    .withWatermark("trade_time", "10 minutes") \
    .groupBy(
        F.window("trade_time", "5 minutes"),
        F.col("ticker")
    ) \
    .agg(
        F.first("price").alias("open"),
        F.max("price").alias("high"),
        F.min("price").alias("low"),
        F.last("price").alias("close"),
        F.sum("volume").alias("volume")
    )

# Write output to Delta Lake with checkpointing
query = ohlc_candles.writeStream \
    .format("delta") \
    .outputMode("append") \
    .option("checkpointLocation", "s3://lake-checkpoints/ohlc-candles/") \
    .start("s3://lake-silver/ohlc_candles_5min/")

query.awaitTermination()
Topic 2.5

Fault Tolerance, Checkpointing & Exactly-Once Semantics

Production Reliability

How does a streaming system recover from unexpected executor crashes without duplicating transactions or dropping records?

WAL & Checkpoint Directory

Spark writes offsets and state updates to Write-Ahead Logs (WAL) in S3. On crash recovery, executors reload the exact committed offsets.

Exactly-Once Semantics (EOS)

Requires 3 conditions: (1) Replayable source (Kafka), (2) Deterministic transformations, and (3) Idempotent or transactional sink (Delta Lake / PostgreSQL).

Dead Letter Queue (DLQ) Routing

Corrupted messages that fail schema decoding are intercepted with try/catch and published to a dedicated DLQ topic with error metadata.

Backpressure Mechanisms

Spark automatically senses downstream bottlenecks and throttles input consumption using PID rate controllers (spark.streaming.backpressure.enabled).

2. Real-World Enterprise Streaming Architecture

Production Pattern
[Exchanges / WebSocket Feeds] ↓ [Kafka Producer] • acks=all • Snappy • Key=Ticker • Schema Registry ↓ [Amazon MSK / Kafka Cluster] (3 AZs • ISR=2 • Log Retention=7 Days) • topic: raw-ticks [p0..p11] • topic: dlq-poison-pills ↓ ↓ [Spark Structured Streaming] [SNS Alert / SQS DLQ] • RocksDB State Store • S3 Checkpoint WAL • 10m Watermark Aggregation ↓ [Delta Lake / S3] • ACID Merge • Real-time Pricing Cache in Redis
Hands-On Lab

Lab 02: Real-Time Trade Stream & Anomaly Detection with Kafka and Spark

Lab Guide
1

Launch Local Kafka Broker

Spin up Kafka using KRaft mode in Docker and create topic financial-trades with 6 partitions.

2

Simulate Event Stream

Run a Python generator publishing 500 simulated stock trade events per second with occasional high-value trade anomalies.

3

Implement Spark Structured Streaming Detector

Subscribe to Kafka, parse JSON, compute rolling 1-minute trade volume z-scores, and flag transactions exceeding 3 standard deviations.

4

Output to Alerts Topic & Sink

Route flagged anomalies to topic fraud-alerts and persist all validated trades to Delta Lake with checkpointing.

3. Streaming & Kafka Interview Preparation

Interview Scenarios
Q1: What causes Kafka consumer rebalance storms and how do you prevent them? ▲

Answer: Rebalance storms occur when consumers in a group repeatedly drop out and trigger partition reassignment. Common causes: (1) max.poll.interval.ms exceeded because message processing took longer than expected; (2) Heartbeat timeouts due to GC pauses; (3) Pod autoscaling thrashing.
Fixes: Increase max.poll.interval.ms, decrease max.poll.records, switch to the CooperativeStickyAssignor (which only pauses migrating partitions rather than revoking all partitions), and optimize GC settings on consumer containers.

Q2: How does Spark Structured Streaming achieve Exactly-Once processing end-to-end? ▼

Answer: Exactly-once requires: (1) Replayable Source: Offsets tracked in Write-Ahead Logs; (2) Deterministic Processing: Pure operations where re-executing batch ID 42 produces identical data; (3) Idempotent or Transactional Sink: The sink must support atomic commits or idempotent upserts (like Delta Lake's transaction log _delta_log or relational upserts). When recovering, Spark discards partial uncommitted writes and replays from the last recorded checkpoint offset.

4. Related Tracks

Explore Next

Continue your journey into Financial Data Engineering

Discover how streaming and batch architectures apply directly to capital markets, ETF pricing, and portfolio NAV calculation.

📈 Explore Financial Data Track →
Open Kafka & streaming interview questions →