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

Stream scoring replay: checkpoints, duplicates and output reconciliation

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

Restoring stream state is not enough; published decisions must also survive source replay without silent duplication.

Name the recovery boundary

A stateful scorer restores feature aggregates from a checkpoint and rereads source events after that position. State may be consistent while an external decision sink already contains some outputs that the restored job emits again. End-to-end exactly-once behavior needs a replayable source and an idempotent or transactional sink. Record the source offsets, checkpoint ID, model digest and output transaction or idempotency key as one recovery packet. Batch idempotency shares the output-identity principle but uses different scheduling boundaries.

Build a stable decision key

A useful key combines immutable event identity and scoring revision. A model redeployment does not automatically authorize a new customer decision for an old event; a deliberate rescore must use a new revision and link to its predecessor. Reject duplicate writes with the same key and same content. If the same key arrives with different model or route data, raise a conflict instead of last-write-wins. Promotion identity makes the model side of that conflict inspectable.

Reconcile after restart

Compare the intended event range with emitted decision keys, sink acknowledgments and dead-letter counts. Gaps may come from a permanently late event, an expired source offset, a failed output transaction or a held schema. Replaying the source should either reproduce the same decision or create a clearly named correction, according to policy. Sampling logs alone cannot prove completeness. Coverage ledgers distinguish sampled diagnostics from required accounting.

Drill the failure path

Kill the worker after state changes but before sink acknowledgment, restore from the prior checkpoint and verify one externally visible decision per key. Repeat with a model pointer change during downtime. The restart should honor the pinned revision for replay or explicitly enter a correction workflow; it must not score old events with an unnamed new model. The project reconciles 47 repeated uploads and a checkpoint rewind.

Implementation

python
def reconcile_decisions(events, published):
    expected = {(event["event_id"], event["scoring_revision"])
                for event in events if event["eligible"]}
    observed = {(row["event_id"], row["scoring_revision"])
                for row in published}
    return {"missing": expected - observed, "unexpected": observed - expected,
            "duplicate_rows": len(published) - len(observed)}

events = [{"event_id": "receipt-47", "scoring_revision": "r8", "eligible": True},
          {"event_id": "receipt-82", "scoring_revision": "r8", "eligible": True}]
published = [{"event_id": "receipt-47", "scoring_revision": "r8"},
             {"event_id": "receipt-47", "scoring_revision": "r8"}]
report = reconcile_decisions(events, published)
assert report["missing"] == {("receipt-82", "r8")}
assert report["duplicate_rows"] == 1

Performance and operating cost

The comparison is O(e + p) expected time and O(e + p) memory for e eligible events and p published rows. Run it in bounded partitions in production. Checkpoint storage and transactional sinks add latency and operational cost; idempotent keys can be simpler when exact transaction coordination is unavailable, provided conflicts and gaps remain observable.

Common Mistakes

  • Calling state checkpointing end-to-end exactly once without checking the sink.
  • Using model version alone as a decision key.
  • Letting replay overwrite a decision with a different route.
  • Measuring only duplicates while missing decisions go unnoticed.

Read next

ai-data
mlops
Storage details