State Stores in Kafka Streams: Local State with Changelog Recovery
Kafka Streams uses state stores because stateful operations like joins, aggregations, and windowing require remembering data across records. A stateless map or filter can process each record independently, but a count per key, a join with a reference table, or a windowed aggregation needs to look up previous values. A state store is a local, durable key-value store that Kafka Streams maintains on the instance running the task. By default, it uses RocksDB, which is an embedded key-value store optimized for SSD and memory. The state store is local to the task, which means lookups are fast and do not require a network round trip. This is a fundamental design choice: Kafka Streams brings the state to the computation rather than sending the computation to a remote state store.
The mechanism that makes this fault-tolerant is the changelog topic. Every state store is backed by a Kafka changelog topic that records every update to the store. When a task is assigned to an instance, Kafka Streams restores the state store from the changelog by replaying it from the beginning, or from a checkpoint if the store already exists locally. If the instance crashes, the task is reassigned to another instance, which restores the state from the changelog and continues processing. This is why Kafka Streams can tolerate failures without a separate storage cluster: the changelog is the durable source of truth, and the local store is a cache that can be rebuilt. The changelog topic is compacted by default, so it retains only the latest value per key, which bounds its size. This is a key trade-off: if the state store is large, restoration from the changelog can take a long time, which is why Kafka Streams also supports standby replicas and local checkpoints.
A common mistake is to assume that the state store is shared across instances. It is not; each task has its own state store, and tasks are partitioned by key. If you need to join two streams, they must be co-partitioned: same number of partitions and same partitioning strategy. If they are not, Kafka Streams will either fail the join or require a repartition, which is expensive. Another mistake is to ignore the restoration time. A large state store with a high update rate can take minutes or hours to restore, during which the task is not processing. This is why monitoring restoration metrics and sizing changelog topics correctly matters. The trade-off is between local state (fast lookups, but restoration cost) and remote state (no restoration, but network latency per lookup). Kafka Streams chooses local state because stream processing is latency-sensitive and network round trips per record would be prohibitive. Version note: RocksDB remains the default, but Kafka Streams also supports in-memory stores for testing and for cases where durability is not needed. KIP-844 (Kafka 3.0) improved restoration with a new mechanism for tracking checkpointing and restoration progress, and standby replicas have been available for a long time as a way to reduce failover time.
State stores hold local key-value state for joins, aggregations, and windowing.
RocksDB is the default; lookups are local and fast with no network round trip.
Every state store is backed by a compacted changelog topic for fault tolerance.
On failover, state is restored by replaying the changelog on the new instance.
Joins require co-partitioned inputs; otherwise a repartition is needed.
Restoration time depends on state size and changelog retention; standby replicas reduce failover time.
KIP-844 (Kafka 3.0) improved restoration; in-memory stores are available for testing.
0-2 years experience
2-5 years experience
5-8 years experience
8+ years experience