Skip to content
Reliable Data Engineering
Lesson
Open in the interactive app

Streaming Deep Dive

Kafka, Spark Structured Streaming, Flink and the streaming patterns that come up in nearly every senior data engineering loop.

Streaming concepts


Kafka Fundamentals

Core Concepts

ConceptDescription
TopicCategory/feed of messages
PartitionOrdered, immutable sequence within a topic
OffsetPosition of a message within a partition
Consumer GroupSet of consumers sharing the work of reading a topic
BrokerKafka server that stores data

Partitioning

Topic: orders (6 partitions)
┌────────────┬────────────┬────────────┐
│ Partition 0│ Partition 1│ Partition 2│
│ user_1     │ user_2     │ user_3     │
│ user_7     │ user_8     │ user_9     │
└────────────┴────────────┴────────────┘
┌────────────┬────────────┬────────────┐
│ Partition 3│ Partition 4│ Partition 5│
│ user_4     │ user_5     │ user_6     │
│ user_10    │ user_11    │ user_12    │
└────────────┴────────────┴────────────┘

Partition Key Choices:

KeyOrdering GuaranteeRisk
user_idAll events for a user in orderHot users cause skew
null (round-robin)NoneEven distribution
event_typeEvents of same type in orderUneven if types vary

Sizing Rules

Partitions = max(
    desired_throughput / throughput_per_partition,
    num_consumers_in_group
)

# Example: 100K msg/sec, each partition handles 10K/sec
# Partitions = max(100K/10K, 20 consumers) = max(10, 20) = 20

Retention

StrategyUse Case
Time-based (7 days)Standard: replay window
Size-based (100GB)When storage constrained
CompactedKeep latest per key (changelogs)

Delivery Semantics

At-Most-Once

At-Least-Once

Exactly-Once

How to Achieve Exactly-Once:

# Kafka producer
producer = KafkaProducer(
    enable_idempotence=True,  # Prevents duplicates on retry
    transactional_id="my-txn-id"  # Enables transactions
)

# Spark consumer
spark.readStream \
    .format("kafka") \
    .option("kafka.isolation.level", "read_committed")  # Only committed txns

Windowing

Tumbling Window

Fixed, non-overlapping intervals.

Events: [1, 2, 3, 4, 5, 6, 7, 8, 9]
Window: 3 events
Result: [1,2,3], [4,5,6], [7,8,9]
.groupBy(window("timestamp", "5 minutes"))

Sliding Window

Overlapping intervals.

Events: [1, 2, 3, 4, 5]
Window: 3, Slide: 1
Result: [1,2,3], [2,3,4], [3,4,5]
.groupBy(window("timestamp", "10 minutes", "2 minutes"))  # 10min window, 2min slide

Session Window

Gap-based, ends after inactivity.

Events: [1, 2, 3, --gap--, 7, 8, --gap--, 15]
Gap: 3
Result: [1,2,3], [7,8], [15]
.groupBy(session_window("timestamp", "30 minutes"))  # Spark 3.2+

Watermarks & Late Data

What is a Watermark?

A watermark says: “I believe I’ve seen all events with timestamp <= W”

df.withWatermark("event_time", "10 minutes")

This means: Allow events up to 10 minutes late. Later events are dropped.

Handling Late Data

StrategyImplementation
Drop late eventsDefault with watermark
Emit late to separate sinkoutputMode("update") + side output
Reprocess in batchDaily reconciliation job
# Late data to separate stream (Flink)
late_data = windowed_stream.getSideOutput(late_output_tag)

Spark Structured Streaming

Basic Pattern

# Read
df = spark.readStream \
    .format("kafka") \
    .option("subscribe", "topic") \
    .option("startingOffsets", "latest") \
    .load()

# Transform
result = df.select(from_json(col("value"), schema).alias("data")) \
    .select("data.*") \
    .withWatermark("timestamp", "5 minutes") \
    .groupBy(window("timestamp", "1 minute")) \
    .count()

# Write
query = result.writeStream \
    .format("delta") \
    .outputMode("append") \
    .option("checkpointLocation", "/checkpoint") \
    .trigger(processingTime="30 seconds") \
    .start()

Output Modes

ModeBehaviorUse Case
AppendOnly new rowsNon-aggregated output
UpdateChanged rows onlyAggregations to KV store
CompleteFull result tableAggregations to file (small)

Triggers

TriggerBehavior
processingTime="10 seconds"Micro-batch every 10s
once=TrueSingle batch, then stop
availableNow=TrueProcess all available, then stop
continuous="1 second"True streaming (experimental)

AspectSpark StreamingFlink
ModelMicro-batchTrue streaming
LatencySeconds-minutesMilliseconds
StateIn-memory + checkpointRocksDB + checkpoint
CEPLimitedFirst-class
Exactly-onceWith checkpointsNative
BatchExcellentGood

When to Use Spark Streaming


State Management

Stateless vs Stateful

TypeExampleComplexity
StatelessFilter, map, projectEasy
StatefulAggregation, join, dedupHard

State Backends

BackendUse Case
In-memory (Spark)Small state, fast
RocksDB (Flink)Large state, persistent

State Cleanup

# Spark: State expires with watermark
df.withWatermark("timestamp", "1 hour")  # State cleared after 1 hour

# Flink: Explicit TTL
descriptor.enableTimeToLive(StateTtlConfig.newBuilder(Time.hours(1)).build())

Checkpointing

Periodic snapshots of state for fault tolerance.

# Spark
.option("checkpointLocation", "/checkpoints/my_job")

# Flink
env.enableCheckpointing(60000)  # Every 60 seconds
env.getCheckpointConfig().setCheckpointStorage("s3://bucket/checkpoints")

Recovery

On failure, restart from last checkpoint. Reprocess events since checkpoint.


Common Patterns

Deduplication

# Spark: Dedupe within watermark window
df.withWatermark("timestamp", "10 minutes") \
  .dropDuplicates(["event_id", "timestamp"])

# Alternative: Use Delta MERGE with event_id as key

Stream-Stream Join

# Join two streams within time range
left.withWatermark("time", "10 minutes") \
    .join(
        right.withWatermark("time", "10 minutes"),
        expr("left.key = right.key AND left.time BETWEEN right.time - interval 5 minutes AND right.time + interval 5 minutes")
    )

Stream-Static Join (Enrichment)

# Join streaming events with static dimension
events.join(broadcast(dim_product), "product_id")

References


Stream-stream join: how state is bounded

sequenceDiagram
    participant I as impressions stream
    participant J as join operator (state store)
    participant C as clicks stream
    I->>J: impression ad=7 t=10:00 (buffered)
    C->>J: click ad=7 t=10:03
    J-->>J: match within 0-15 min interval
    Note over J: emits joined row
    Note over J: watermark passes 10:15 + delay, so impression state is evicted
imps   = impressions.withWatermark("imp_ts", "10 minutes")
clicks = clicks.withWatermark("click_ts", "20 minutes")
joined = imps.join(
    clicks,
    expr("""click_ad_id = imp_ad_id AND
            click_ts >= imp_ts AND
            click_ts <= imp_ts + interval 15 minutes"""),
    "leftOuter")          # outer joins REQUIRE watermarks + time bounds

Without watermarks and a time-range condition, the engine must keep both sides forever, and state grows without bound.


Interview questions

Your streaming job's state store keeps growing until executors OOM. What do you check?
  1. Is there a watermark on every stateful operator (aggregations, dedup, stream-stream joins)? 2. Do joins have a time-range bound? 3. Is the key cardinality unbounded (e.g. grouping by session_id with no window)? 4. Use RocksDB state store instead of in-memory (HDFS-backed) for large state. 5. For arbitrary stateful processing, set state TTL / timeouts. 6. Watch the stateOperators.numRowsTotal and memory metrics in the streaming query progress.
How do you deduplicate a stream exactly?

Pick a unique event id. dropDuplicatesWithinWatermark("event_id") with a watermark bounded by the max expected duplicate delay; plus an idempotent sink (MERGE on event_id) to catch anything older than the watermark. Without a watermark, dropDuplicates keeps every id forever.

Checkpoint vs savepoint vs Delta transaction log?

Checkpoint: automatic, engine-owned snapshot of offsets + state for failure recovery. Savepoint (Flink): user-triggered, portable snapshot for upgrades/migrations/rescaling. Delta log: the sink table’s own commit history; Structured Streaming also writes (appId, batchId) there so replays don’t double-write.

Kafka consumer lag is growing steadily. Walk through your diagnosis.

Is input rate up (traffic spike) or processing rate down (slower batches)? Check batch duration vs trigger interval, skewed partitions (one partition lagging = hot key), GC/OOM, slow sinks (DB backpressure), external lookups per record. Fixes: scale executors / partitions (if parallelism is capped by partition count, add partitions carefully), salt hot keys, batch external calls, tune maxOffsetsPerTrigger, optimise the sink.