skip to content

Where does stateful windowing state live in Kafka Streams versus Flink versus Spark, and how is it made fault-tolerant?

level: middleimportance: should knowfreq 55%

answer

  1. Streams: RocksDB + compacted changelog topic
  2. standby replicas cut restore time
  3. Flink: state backend + barrier checkpoints to S3
  4. savepoints for upgrades/rescale
  5. Spark: checkpoint dir = offsets + versioned state

basics

~20 s

Kafka Streams keeps state in embedded RocksDB on each instance, backed up to compacted changelog topics in Kafka. Flink keeps state in a state backend (heap or RocksDB) snapshotted to durable storage via checkpoints. Spark stores state in a checkpointed, versioned state store.

solid answer

~40 s

Stateful operations (aggregations, joins, windows) need durable, recoverable state. Kafka Streams stores it in per-instance embedded RocksDB and continuously mirrors every change to a compacted changelog topic in Kafka; on failure or rebalance, a new instance restores the store by replaying that changelog (standby replicas keep warm copies to cut restore time). Flink stores state in a configurable state backend — HashMapStateBackend (on-heap) or EmbeddedRocksDBStateBackend (off-heap, supports large state and incremental checkpoints) — and periodically snapshots it to durable storage (S3/HDFS) via asynchronous barrier checkpointing; recovery restores from the latest checkpoint or an operator-triggered savepoint. Spark Structured Streaming uses a checkpoint location holding offset/commit logs plus a versioned state store (HDFS-backed or RocksDB), recovering by replaying from the last committed batch. The recurring theme: local fast state plus durable remote backup.

go deeper

for a junior

Know state is stored locally and backed up durably so it can recover after a crash.

for a middle

Name the mechanisms: Streams RocksDB + changelog, Flink state backend + checkpoints, Spark checkpoint dir + state store.

for a senior

Discuss restore time, standby replicas, incremental checkpoints, and savepoints for upgrades/rescaling.

for a principal

Weigh state-size limits, restore SLAs, and durable-store dependencies when choosing an engine for large stateful workloads.

## Why state needs special handling Stateless transforms (map, filter) can restart anywhere. **Stateful** ones — counting per key, windowed aggregates, joins, deduplication — accumulate data that must survive crashes and must move correctly when work is rebalanced. Each engine pairs **fast local state** with a **durable backup**. ## Kafka Streams - **Local store**: each task holds an embedded **RocksDB** instance on local disk (in-memory stores are also available for small state). Lookups and updates are local and fast. - **Durability**: every write to a state store is also appended to a **changelog topic** in Kafka, which is **log-compacted** so it retains the latest value per key. The changelog is the source of truth. - **Recovery / rebalance**: when an instance dies or partitions move, the receiving instance **restores** the store by consuming its changelog from the beginning (compacted, so bounded by key count). To avoid slow restores, configure **standby replicas** (`num.standby.replicas`) that keep warm copies; Streams can also use **interactive queries** to read local state directly. - **Operational note**: state is tied to disk on the app instances. Losing a node means restoring from changelog; large state means long restores without standbys. ## Apache Flink - **State backends**: `HashMapStateBackend` keeps state as objects **on the JVM heap** (fast, limited by memory); `EmbeddedRocksDBStateBackend` keeps it **off-heap in RocksDB** on local disk, enabling **very large state** and **incremental checkpoints** (only changed RocksDB SST files are uploaded). - **Checkpointing**: Flink periodically runs **asynchronous barrier snapshots** — barriers flow through the dataflow and each operator snapshots its state to **durable storage** (S3, HDFS, etc.). Checkpoints are for automatic failure recovery. - **Savepoints**: user-triggered, portable snapshots used for upgrades, rescaling, and code changes. They are the operational lever for stateful redeploys. - **Recovery**: on failure Flink restores all operator state from the last completed checkpoint and rewinds sources to the matching offsets. ## Spark Structured Streaming - **Checkpoint location**: a directory on durable storage containing the **offset log** (which input ranges per batch), the **commit log**, and the **state store** data. - **State store**: versioned per batch; backed by HDFS-compatible storage by default, with an optional **RocksDB state store** for large state to reduce JVM GC pressure. - **Recovery**: on restart Spark reads the checkpoint, identifies the last committed batch, and resumes; state is restored from the versioned store. Watermarks bound how much windowed/join state is retained. ## Comparison | Aspect | Kafka Streams | Flink | Spark | |---|---|---|---| | Local state | RocksDB / in-memory | heap or RocksDB | heap or RocksDB | | Durable backup | compacted changelog topic (in Kafka) | checkpoints to S3/HDFS | checkpoint dir (offset+state) | | Incremental backup | n/a (changelog is per-write) | yes (RocksDB incremental) | versioned state files | | Manual snapshot for upgrades | n/a (changelog replay) | savepoints | restart from checkpoint | ## Edge cases - **Restore time**: huge Streams state without standbys means long rebalance pauses; size standbys vs cost. - **Checkpoint storage failures** (Flink/Spark) can stall the job; the durable store is a hard dependency. - **Schema/topology changes**: Streams changelog and Flink savepoints both have compatibility rules; changing key/serde or operator UIDs can break restore.

  • Why are Kafka Streams changelog topics log-compacted rather than time-retained?
    Because the store only needs the latest value per key to be reconstructed. Compaction keeps the most recent record per key and discards superseded ones, so restore time and storage scale with the key cardinality, not the full update history.
  • What is the difference between a Flink checkpoint and a savepoint?
    Checkpoints are automatic, periodic, owned by Flink for failure recovery, and may be incremental/cleaned up. Savepoints are user-triggered, self-contained, portable snapshots used for planned operations like upgrades, code changes, and rescaling.

saying these in an interview costs you the question

  • Saying Kafka Streams state lives 'in the Kafka cluster' — it lives in local RocksDB; the changelog is the backup in Kafka.
  • Claiming Flink only keeps state on-heap — RocksDB backend supports very large off-heap state with incremental checkpoints.
  • Assuming any engine survives loss of its durable backup store; checkpoints/changelogs are a hard dependency.
  • Conflating checkpoints and savepoints in Flink.

context