A stateful stream operator remembers values by key; a checkpoint binds that state to an input position so recovery can replay consistently.
Keyed state, checkpoints and recovery
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
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_statePerformance 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
- Event time and late arrivals: close windows with an explicit correction policy
- Idempotent loads: commit target rows and extraction progress together
- Stream backpressure and catch-up budget
- Project: release a recoverable streaming risk ledger
- Idempotent streaming sinks: make retries safe across failures
Continue the workflow: Failover writer fencing and source positions.
