State Stores

State Stores, Changelogs and Fault Tolerance

Each store is a RocksDB database per task under state.dir, and every write also goes to a compacted changelog topic, so the files are only a cache: a task that moves to another instance is rebuilt from its changelog. After the 584 KB state directory here was deleted, a fresh start rebuilt both stores and settled after 53 seconds.

Internal topics carry the application.id as a prefix, so never change it in production. A restart with intact files replays only what it missed (since 4.3 each store tracks its changelog offsets, KIP-1035), and standby replicas keep warm copies elsewhere.