State in a streaming job represents the intermediate data maintained across records to compute aggregations, windowed joins, and deduplication logic. This state is stored in memory or embedded key-value engines like RocksDB and is periodically written to durable checkpoints for fault tolerance. State grows uncontrollably when pipelines process high-cardinality keys without advancing watermarks or state time-to-live expiration policies, leading to prolonged checkpoint intervals and worker out-of-memory crashes.
Where state accumulates in operators
Stateful stream processing requires operators to remember historical information across distinct input events:
- Windowed aggregations: Operators accumulate counters, running sums, or lists of records per key until the window boundary closes.
- Stream-to-stream joins: Both sides of an inner or outer join must buffer records in state while waiting for matching timestamps from the opposite stream.
- Deduplication filters: Operators retain previously observed record identifiers to prevent duplicate processing on stream retries.
In frameworks like Apache Flink or Spark Structured Streaming, this keyed state is stored in RocksDB instances on worker local disks or within JVM heap memory.
Root causes of unbounded state growth
The primary driver of state explosion is processing unbounded high-cardinality keys without eviction rules. Tracking metrics by unique user IDs, ephemeral browser session IDs, or transaction UUIDs generates millions of distinct keys every day.
If the job lacks an advancing watermark to finalize windows, or fails to define a State Time-To-Live (TTL), the storage backend retains every historical key indefinitely. Over days of continuous execution, the local RocksDB data directory swells to hundreds of gigabytes.
Production symptoms and remediation
When state size becomes unmanageable, several operational symptoms emerge:
- Checkpoint durations balloon from seconds to tens of minutes because workers must serialize and transfer massive state snapshots to cloud object storage.
- Worker tasks suffer severe Garbage Collection pauses or fail completely with out-of-memory errors during state materialization.
To resolve state bloat, apply three corrective measures:
- Configure explicit state TTL rules that automatically purge keys unaccessed after a specified duration, such as twenty-four hours.
- Implement strict watermarks on event-time windows to trigger prompt state eviction upon window closure.
- Reduce key footprint by hashing string identifiers into compact integer keys or storing minimal accumulator objects instead of raw event payloads.