Flink
Mid-levelHow do I configure Flink's state backend and checkpoint storage before going to production?
Flink's defaults are tuned for getting a job running fast in development, not for surviving a real restart with real state size. Two settings decide almost everything.
Every stateful Flink job needs a state backend (where operator state lives while the job runs) and checkpoint storage (where consistent snapshots of that state get written). Flink runs fine with the defaults right up until state grows past what a JVM heap comfortably holds.
HashMapStateBackend vs EmbeddedRocksDBStateBackend
HashMapStateBackend keeps all state as Java objects on the JVM heap - fast, but bounded by
however much heap you've given the task manager, and every checkpoint has to serialize the entire state
being written. EmbeddedRocksDBStateBackend keeps state on local disk (with a configurable
in-memory cache) and only serializes what's actually needed per checkpoint, at the cost of some latency per
state access. Past a few hundred MB of state per task slot, RocksDB is usually the right default, not an
optimization to reach for later.
state.backend.incremental: true matters almost as much as the backend choice itself - without
it, every checkpoint re-writes the full state snapshot instead of just what changed since the last one,
which gets expensive fast as state grows.
Checkpoint storage has to be durable and shared, not local
Checkpoint storage is a separate setting from the state backend - it's where the actual checkpoint files land, and it has to be reachable by whichever task manager recovers a failed task, which in practice means object storage, not local disk:
num-retained keeps more than one checkpoint around on purpose - if the most recent checkpoint
turns out to be corrupt or mid-write when a task manager dies, Flink falls back to the previous one instead
of having nothing to recover from.
The setting that's painful to change later
Switching state backends after a job has real production state means the new backend can't read the old backend's checkpoint format - there's no in-place migration, only a savepoint taken under the old backend and restored under the new one, with real downtime while that happens. Picking RocksDB before launch if there's any chance state will grow costs nothing; migrating off HashMapStateBackend under load is a planned maintenance window.
Outside opti-pipe's scope: state backend and checkpoint storage are cluster/job configuration, set before a job starts - the rule engine only reads the metrics a running or completed job exports, it doesn't touch how or where state gets persisted.
See what this looks like on your own pipeline.
Upload a real Spark event log, dbt run_results.json, or Flink metrics export and get
concrete, approve-before-apply recommendations back - not another rule of thumb.