Stateful Streaming — Keyed State, RocksDB, Checkpoints & Rescaling
Stateful streaming means the operator remembers things between events: running counts, open windows, join buffers, dedupe sets, feature rows. That state can be terabytes. It must survive crashes consistently with the source offsets, and it must be redistributable when you change parallelism. Flink's reference design is keyed state partitioned into key groups, stored in RocksDB, snapshotted by barrier checkpoints into object storage, and restorable through savepoints. Spark and Kafka Streams solve the same problem with different trade-offs.
- 1Gist
- 2Maps
- 3Q&A
- 4Sandbox
Voice readout needs Web Speech Synthesis in this browser.
Question ladder
L1
How does Flink take a consistent snapshot without pausing the job?
Answer
Barriers are injected at sources and flow with the data. Each operator snapshots when it has seen the barrier on all inputs. The snapshot is asynchronous. The union of snapshots plus source offsets is a consistent cut.
L2
What is the difference between aligned and unaligned checkpoints?
Answer
Aligned waits for the barrier on every input and buffers faster inputs. Snapshots stay small and slow down under backpressure. Unaligned lets the barrier overtake in-flight data and persists that data, so checkpoints finish under backpressure and are larger.
L3
Why RocksDB instead of heap state?
Answer
State bigger than memory, incremental checkpoints, and predictable garbage collection. The trade-off is serialization and disk latency on every access.
L4
What is a checkpoint versus a savepoint?
Answer
Checkpoints are engine-managed for recovery and may be deleted. Savepoints are user-triggered portable snapshots for upgrades, rescaling, and migrations. They need stable operator UIDs.
L5
How is keyed state redistributed when you scale from 4 to 6?
Answer
Keys map to a fixed number of key groups, maxParallelism. Each subtask owns a contiguous range. Rescaling reassigns ranges and restores them from the last snapshot.
L6
How do you keep state bounded?
Answer
Windows with allowed lateness, watermark-driven cleanup in joins, state TTL for dedupe and features, and timers that delete keys. Monitor state size per operator.
L7
How do Kafka Streams and Flink differ on state durability?
Answer
Kafka Streams writes every state change to a compacted changelog and restores by replaying it. Flink snapshots periodically to object storage and restores by downloading the snapshot plus replaying the source from the saved offsets.
Failure modes
Offsets committed apart from state
Auto-commit moves the offset. The snapshot is older. The records in between are lost. The reverse order duplicates.
No operator UIDs
A refactor changes generated ids and the savepoint cannot restore.
maxParallelism too low
You cannot add parallelism in a spike without dropping keyed state.
Unbounded dedupe or join state
RocksDB grows until the disk fills.
Misconceptions
A checkpoint pauses the job.
Alignment buffers fast inputs. The snapshot itself is asynchronous. Processing continues while RocksDB uploads.
Key groups avoid moving keys.
Going from 2 to 3 subtasks still moves about half the keys. They move as contiguous ranges, which is what makes restore cheap.
Hash of the key modulo parallelism is the same thing.
That scatters every key on every rescale. Key groups are fixed for the life of the job.
Interviewer traps
Putting a giant list in one ValueState.
Every update rewrites the blob. Use MapState or ListState so you touch only what changed.
Treating a Kafka Streams restore as a Flink checkpoint download.
Changelog replay can take minutes. Standby replicas exist to cut that time.
Design scenario
Same prompt for every reader.
Requirements
RocksDB, aligned or unaligned checkpoints to object storage, key groups with a generous maxParallelism, stable UIDs, and TTL on the dedupe set.
Traffic / scale
A few hundred thousand events per second, state in the hundreds of gigabytes.
Latency
Checkpoint interval of about one minute. Restore should replay only the tail after the last completed checkpoint.
Consistency
Source offsets and keyed state are one cut. A naive auto-commit must fail the interview.
Availability
Unaligned checkpoints if backpressure would otherwise stall alignment. Savepoint before the rescale.
Failure assumptions
- The job dies after a barrier is snapshotted and before the next one.
- Parallelism changes from 4 to 6.
Constraints
- Do not change maxParallelism after the first deploy.
- Do not store the dedupe set as one blob per key.
Prompt
A dedupe set and a five-minute window count share one Flink job. You need to double parallelism for a sale, and a crash must not lose counts that were already offset-committed.
API
Which operator UID must stay stable across the sale deploy?
Data
What is in the checkpoint besides the counts?
Architecture
Where do the SST files and the offsets live, and how does a key group move?
One cut, or two clocks
Prefer
Offsets and state at the same barrier
Restore rewinds the source to the snapshotted offset and reloads the counts from that cut. Replay finishes the tail once.
- Records before the barrier belong to this checkpoint.
- The snapshot uploads in the background.
- Transactional sinks commit only after every task acks.
Alternative
The consumer auto-commits, state saves later
The offset is new and the counts are old. The records in between never land in state. Flip the order and you duplicate instead.
- A crash looks like data loss with no error in the log.
- A second crash looks like double-counting.
- No rescale story, because keys were hash modulo parallelism.
A barrier becomes a checkpoint
Chandy-Lamport, adapted to a streaming DAG. The job does not stop the world.
- 1
Coordinator triggers checkpoint n
Sources record their partition offsets as the barrier is injected. - 2
Barrier flows in-band
Records before the barrier belong to n. Records after it belong to n+1. - 3
Operator aligns inputs
It waits for barrier n on every input, buffering faster ones, then snapshots. - 4
Async upload, then ack
RocksDB snapshots locally and uploads while processing continues. The checkpoint completes when every task acks.
Kinds of state
| State | Scope | Examples | Redistribution |
|---|---|---|---|
| Keyed | Per key, after keyBy | A count, a cart map, a buffer list, window contents, timers | By key-group ranges |
| Operator, non-keyed | Per parallel instance | Kafka source offsets, sink transaction handles | Even split, or union |
| Broadcast | The same copy on every instance | Rules, flags, config pushed as a stream | Copied |
Keyed state is the one that grows. Bound it with TTL or with watermark-driven cleanup from the watermark lesson.
Backends
| Backend | Where live state sits | Snapshot | Practical size | Latency | Pick when |
|---|---|---|---|---|---|
| Heap | JVM objects | Full copy each checkpoint | Fits in memory, gigabytes | Fastest | Small state, low latency |
| RocksDB | Local SSD, an LSM, off-heap | Incremental: upload new SST files | Terabytes per job | Slower. Serialization plus disk | Large windows and joins. The production default |
| Disaggregated, ForSt in Flink 2.x | Remote storage, local cache | Native to remote storage | Very large, fast rescale | Higher, mitigated by an async state API | Elastic cloud jobs |
| Spark state store | HDFS-backed map or RocksDB | Versioned per micro-batch in the checkpoint dir | Medium to large | Per batch | Structured Streaming stateful ops |
| Kafka Streams | Local RocksDB per task | A compacted changelog replicates every write | Large. Restore replays the changelog | Fast reads | Kafka-native apps. Standby replicas cut restore time |
RocksDB makes state size a disk problem and enables incremental checkpoints. Every access serializes. Do not put a giant list in one value. Use map or list state so an update touches one entry. Snapshots in object storage follow object storage. Restore time after a rebalance is the Kafka-side pain in consumer lag.
Barriers
Sequence
- 1
Coordinator → Source
Inject barrier n
- 2
Source → Operator
Barrier n plus offsets
- 3
Operator → State backend
Snapshot state
- 4
Operator → Sink
Pre-commit
- 5
Source → Coordinator
Ack
- 6
Operator → Coordinator
Ack
- 7
Sink → Coordinator
Ack
- 8
Coordinator → Sink
Checkpoint complete, commit
Lesson map
Stateful Streaming — Keyed State, RocksDB, Checkpoints & Rescaling
Stateful streaming means the operator remembers things between events: running counts, open windows, join buffers, dedupe sets, feature rows. That state can be terabytes. It must survive crashes consistently with the source offsets, and it must be redistributable when you change parallelism. Flink's reference design is keyed state partitioned into key groups, stored in RocksDB, snapshotted by barrier checkpoints into object storage, and restorable through savepoints. Spark and Kafka Streams solve the same problem with different trade-offs.
Architecture. Architecture
Select a node to see why it exists, or an edge to see the protocol, direction, effect, and consequence.
Mermaid export
flowchart TB jm["Coordinator"] src["Source"] op["Operator"] jm -->|Trigger| src src -->|Barrier n with| op op -->|Ack after state| jm jm -->|Checkpoint n| op
- The coordinator injects barrier n at the sources, which record their offsets.
- Barriers flow with the data. Records before the barrier belong to checkpoint n. Records after it belong to n+1.
- An operator with several inputs aligns: it waits for barrier n on every input, buffering faster inputs, then snapshots. That is a consistent cut without stopping the world.
- Snapshots are asynchronous. RocksDB takes a local snapshot and uploads in the background while processing continues.
- When every task acks, checkpoint n is complete. Transactional sinks commit. That commit is the exactly-once lesson.
Unaligned checkpoints, since Flink 1.11, let barriers overtake buffered records and store in-flight data in the checkpoint. They complete under backpressure and the snapshots are bigger. Buffer debloating reduces in-flight data so aligned checkpoints stay fast.
Press Run. Snippets must be self-contained — no network, files, or native modules.
truth {'a': 5, 'b': 3, 'c': 2}
aligned barrier {'a': 5, 'b': 3, 'c': 2} == truth: True
naive (auto-ack) {'a': 4, 'b': 3, 'c': 1} == truth: FalseThe naive job lost one a and one c. Those records were committed at the source and their effect was not in the saved snapshot. Save state first and commit offsets later, and you get duplicates. Only one consistent cut gives exactly-once state.
Checkpoints versus savepoints
| Checkpoint | Savepoint | |
|---|---|---|
| Purpose | Automatic failure recovery | Planned work: upgrade code, change parallelism, migrate, try a new version |
| Owner | The engine. It may delete old ones | You. Never auto-deleted |
| Format | Backend-native, may be incremental | Canonical or native, self-contained |
| Frequency | Every N seconds or minutes | On demand |
| Requirement | None beyond a working backend | Stable operator UIDs so state maps onto the new job graph |
Key groups
Keyed state is not hash(key) % parallelism. That would scatter every key on rescale. Instead keyGroup = hash(key) % maxParallelism is fixed for the life of the job, and each subtask owns a contiguous range. Rescaling reassigns whole ranges, so restore reads contiguous chunks.
Press Run. Snippets must be self-contained — no network, files, or native modules.
rescale 2 -> 3 subtasks over 1000 keys
naive hash%parallelism : 671 keys change owner (random reads of scattered state)
key groups (max=128) : 494 keys change owner, moved as 63 whole key groups (read as contiguous ranges, no per-key rehash)
parallelism 2: subtask0=[0..63] subtask1=[64..127]
parallelism 3: subtask0=[0..42] subtask1=[43..85] subtask2=[86..127]
state TTL 3600s at t=7200: live=user-2,user-3 expired=1Key groups do not magically avoid moving keys. About half still change owner from 2 to 3 subtasks. They move as pre-partitioned ranges, which is what makes restore and rescale efficient. maxParallelism caps how far you can ever scale and cannot change without losing keyed state. Set it generously on day one, 720 or 1024 or more.
Keeping state bounded
- State TTL expires an entry N after the last write or read. Cleanup is lazy on access, plus incremental cleanup, or a RocksDB compaction filter.
- Timers fire at window end plus lateness and delete state on purpose.
- Joins and windows drop state when the watermark proves a row can never match.
- In Spark, aggregation state is evicted only when a watermark is set. Arbitrary stateful operators need an explicit timeout or TTL.
Pitfalls
In the Python, point at the records the auto-commit path loses. Then say what changes if you snapshot state on the barrier and only afterwards commit the offset.
Interview Q&A
How does Flink take a consistent snapshot without pausing the job?
Answer
Barriers are injected at sources and flow with the data. Each operator snapshots when it has seen the barrier on all inputs. The snapshot is asynchronous. The union of snapshots plus source offsets is a consistent cut, Chandy-Lamport adapted for a DAG.
Aligned versus unaligned checkpoints?
Answer
Aligned waits for barriers on all inputs and buffers faster inputs. Snapshots stay small and get slow under backpressure. Unaligned lets barriers overtake in-flight data and persists that data. Checkpoints finish under backpressure and are larger.
Why RocksDB over heap state?
Answer
State bigger than memory, incremental checkpoints, and predictable garbage collection. The trade-off is serialization cost and disk latency per access.
Checkpoint versus savepoint?
Answer
Checkpoints are engine-managed for recovery. Savepoints are user-triggered portable snapshots for upgrades, rescaling, and migrations. They rely on stable operator UIDs.
How is keyed state redistributed when you scale from 4 to 6?
Answer
Keys map to a fixed number of key groups, maxParallelism. Each subtask owns a contiguous range. Rescaling reassigns ranges and restores them from the last snapshot.
How do you keep state bounded?
Answer
Windows with allowed lateness, watermark-driven cleanup in joins, state TTL for dedupe and features, and timers that delete keys. Monitor state size per operator.
How do Kafka Streams and Flink differ on state durability?
Answer
Kafka Streams writes every state change to a compacted changelog topic. Restore replays that log. Flink snapshots periodically to object storage. Restore downloads the snapshot and replays the source from the saved offsets.