A corrected historical batch is a new decision version; consumers need a clear supersession rule and reconciliation report.
Batch replay: supersede outputs without duplicating downstream actions
Name the purpose of a replay
A backfill may reconstruct missing scores, apply a corrected feature value or compare a new model. Those are different operations. A comparison run should not replace production decisions. A correction run may supersede earlier outputs, but only after approval and consumer readiness. Keep original and replacement job keys, the changed input or code, and the effective publication time. Batch identity makes two runs for the same calendar partition distinguishable.
Protect downstream side effects
A scoring table can be recomputed; a customer email or manual-review assignment may not be safely repeated. Publish outputs to an immutable versioned location and let consumers advance a pointer or process a change feed keyed by decision version. Require idempotency keys for any action. A replay must state whether it is historical analysis or an operational correction. Never send both old and new decisions to an append-only action queue without a supersession protocol.
Reconcile old and new populations
Count unchanged IDs, added IDs, removed IDs and changed scores. Confirm every input has one scored, rejected or held outcome. If a source correction removes a receipt, do not leave its old score silently active. Record tombstones or explicit withdrawal events where the consumer contract requires them. Outcome joins should preserve which decision version an outcome evaluates; a replayed score is not a retroactive production prediction.
Gate publication with a diff
Review a sample of changed decisions and all high-impact changes before switching the active pointer. Simulate a downstream consumer that has not acknowledged the new version and hold publication. The batch project checks record accounting, supersession and no duplicate customer action under replay. Store the diff manifest for audit and rollback.
Implementation
def reconcile_batch(previous, replacement):
old_ids, new_ids = set(previous), set(replacement)
shared = old_ids & new_ids
changed = {receipt_id for receipt_id in shared
if previous[receipt_id] != replacement[receipt_id]}
return {"added": sorted(new_ids - old_ids),
"removed": sorted(old_ids - new_ids),
"changed": sorted(changed),
"unchanged_count": len(shared - changed)}
old = {"receipt-47": "review", "receipt-82": "clear"}
new = {"receipt-47": "clear", "receipt-129": "review"}
assert reconcile_batch(old, new) == {
"added": ["receipt-129"], "removed": ["receipt-82"],
"changed": ["receipt-47"], "unchanged_count": 0}
Performance and operating cost
Diffing two maps with n and m IDs takes O(n + m) expected time and O(n + m) space; sorting output lists adds O(k log k) for k changed IDs. Keeping versions costs storage, but it is safer than overwriting a production output with no way to explain changed downstream actions.
Common Mistakes
- Treating a model-comparison replay as an approved operational correction.
- Sending replayed decisions through the original action queue.
- Ignoring receipts removed by a corrected source snapshot.
- Joining an outcome to a replayed score instead of its original production decision.
Read next
- Batch inference: make partitions idempotent and outputs identifiable
- Project: publish and replay a nightly receipt scoring partition
- Prediction-outcome joins: evaluate only mature, matched decisions
- Label corrections: version outcomes before rebuilding quality metrics
- Training manifests: link data, code, configuration and artifact
