Skip to content
AITroveRead. Build. Understand.
Make this comfortable

Keyed state, checkpoints and recovery

Last updated: 6 Oct 20265 min read
tutorial
AdvancedBy AITrove Editorial

A stateful stream operator remembers values by key; a checkpoint binds that state to an input position so recovery can replay consistently.

Define the state boundary

A payment-risk counter is keyed by account, not by worker. The engine may move account 47 to another worker after scaling, so the counter must be restorable independently of process memory. Record the input offsets, keyed state and operator version as one recoverable point. Batch checkpoints follow the same replay principle, although stream state persists across many windows.

Treat a checkpoint as a consistent cut

A checkpoint cannot pair state that has processed event 83 with a source offset that resumes at event 79; replay would count four events twice. In-flight records and channel alignment matter when operators receive multiple inputs. An unaligned mode can store those in-flight buffers, trading checkpoint speed under backpressure for larger snapshots. Measure completion time and snapshot bytes together.

Separate engine recovery from sink effects

A restored counter may be exact while an external notification sent just before failure is repeated. Use a transactional sink tied to checkpoint completion or an idempotency key derived from business event identity. Idempotent sinks explain that boundary. Do not promise end-to-end exactly once from an engine checkpoint alone.

Plan state growth and migration

A key that never expires remains in a keyed store indefinitely. Set time-to-live from business retention and late-event policy, then document what happens when a dormant key returns. Serializer or schema changes can make an old checkpoint unreadable. Rehearse restore from a retained snapshot before the release, including a rescale to a different worker count.

Exercise failure between checkpoints

Feed a pinned event log, capture a checkpoint after event 47, process event 48, then stop the worker before the next checkpoint. Recover from event 47 and assert the published output has one effect for event 48. Also test corrupted state, a missing checkpoint file and an incompatible serializer; each needs an explicit stop or rollback path.

Implementation

python
events = [("acct-47", "evt-47", 1200), ("acct-47", "evt-48", 900)]

def apply_events(state, seen, records):
    totals = dict(state)
    identities = set(seen)
    for account_id, event_id, cents in records:
        if event_id in identities:
            continue
        totals[account_id] = totals.get(account_id, 0) + cents
        identities.add(event_id)
    return totals, identities

checkpoint_state, checkpoint_seen = apply_events({}, set(), events[:1])
recovered_state, recovered_seen = apply_events(checkpoint_state, checkpoint_seen, events[1:])
assert recovered_state == {"acct-47": 2100}
assert apply_events(recovered_state, recovered_seen, events[1:])[0] == recovered_state

Performance and operating cost

Processing N records takes O(N) expected time; keyed totals occupy O(K) memory for K active accounts and deduplication occupies O(E) for retained event identities. A real engine stores this state outside the Python process. Checkpoint I/O and restore duration grow with state size and buffered records, so retention and snapshot cadence are operating decisions.

Common Mistakes

  • Do not restart from a source offset newer than the restored state.
  • Do not call an external side effect exactly once without a sink protocol.
  • Do not change state serialization without a restore rehearsal.

Read next

Continue the workflow: Failover writer fencing and source positions.

ai-data
data-engineering
Storage details