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 ShiftTraditional 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.
Apache Kafka Fundamentals & Cluster Architecture
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).
Producers, Consumer Groups & Rebalance Protocol
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.
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()
Schema Management: Schema Registry & Avro/Protobuf
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%.
Spark Structured Streaming & Watermarking
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.
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()
Fault Tolerance, Checkpointing & Exactly-Once Semantics
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 PatternLab 02: Real-Time Trade Stream & Anomaly Detection with Kafka and Spark
Launch Local Kafka Broker
Spin up Kafka using KRaft mode in Docker and create topic financial-trades with 6 partitions.
Simulate Event Stream
Run a Python generator publishing 500 simulated stock trade events per second with occasional high-value trade anomalies.
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.
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
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.
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📈 Financial Data
Real-time tick feeds, ETF pricing, and high-frequency market data pipelines.
Go to Financial Data →⚡ Engineering
Python fundamentals, SQL window functions, and distributed PySpark tuning.
Go to Engineering →☁️ Cloud
Amazon MSK setup, S3 Lakehouse architecture, and AWS MWAA.
Go to Cloud →⚙️ Orchestration
Airflow sensor triggers, streaming health monitors, and CI/CD pipelines.
Go to Orchestration →