advanced 3 min answer Multiple choice

A keyed aggregation holds about 2 TB of state across 48 task slots on a heap state backend with 64 GB of heap per task manager. Checkpoints take 11 minutes and the job dies with multi-second garbage-collection pauses about once a day. Throughput is within target. Which change addresses the actual constraint?

flinkstate-backendrocksdbcheckpointingrecovery
Pick one
Show the full answer Hide the answer

The deciding property

Not throughput. The working set is larger than a managed heap can hold without the collector becoming the bottleneck. 2 TB over 48 slots is about 42 GB of live state per slot against 64 GB of heap. A tracing collector's work scales with the live set, so at that occupancy pauses grow until a heartbeat is missed and the task is declared dead.

A heap backend keeps state as live Java objects: no serialisation on access, and every byte is something the collector must walk. Its snapshot is a full serialisation of everything, which is where the 11 minutes comes from.

An embedded key-value backend keeps state off-heap in a log-structured store on local disk. Capacity is bounded by disk rather than heap, so 42 GB per slot on local NVMe is unremarkable, and JVM memory becomes a tunable block-cache and write-buffer budget of a few GB per slot.

What it costs is per-record: every read and write becomes a serialise plus deserialise and sometimes a disk read. For access-heavy operators that is an order-of-magnitude increase in per-state-access cost. The mitigations are fewer and larger state accesses, a block cache big enough for the hot key set, and a state layout with one lookup per record rather than several.

The recovery arithmetic

Incremental checkpoints upload only files new since the last checkpoint, so checkpoint duration collapses from 11 minutes to the size of the delta and a short interval becomes affordable again.

Recovery does not improve proportionally, and can get worse. Flink's own state-backend documentation says restoring from an incremental checkpoint may take longer when network bandwidth is the bottleneck, because more files have to be fetched - the chain of deltas rather than one snapshot. Restoring 2 TB at an aggregate 1 GB/s is still roughly 35 minutes before the first record moves. The levers that actually cut recovery are task-local state so a restart on the same node reads from local disk, a standby task manager so the rescheduling wait disappears, and smaller state per slot.

So publish two numbers, not one: checkpoint duration and measured restore time.

Why the other options fail

  • Double the heap. It buys one doubling against a state size that is growing, and a larger live set makes pauses longer, not shorter. Raising the checkpoint interval to 10 minutes also means up to 10 minutes of reprocessing after every failure.
  • Shorten the interval. With full snapshots, more frequent checkpoints do not carry less state
  • each one still serialises everything. At 11 minutes per checkpoint and an interval under that, checkpoints overlap or queue, and the job spends most of its time snapshotting.
  • Add 48 slots. Halving per-slot state does postpone the pauses, at double the cluster bill, and leaves the full-snapshot cost untouched. It also requires a restart from a savepoint with key-group redistribution, so it is not the cheap change it looks like.

When this is the wrong answer

If the 2 TB exists because nothing expires - a missing state time-to-live, an unbounded join window, a key space that includes ids never seen again - then a new backend buys you a few months of the same curve at higher per-record cost. Bound the state first, choose the backend second. The tell is state size growing while the active key count is flat.