Build a payment-risk ledger that joins historical merchant categories, revises late windows and recovers without duplicate published effects.
Project: release a recoverable streaming risk ledger
Write the consumer contract
Each payment has a globally unique event ID, account ID, merchant ID, integer cents, event time and arrival time. Publish one account-hour total and risk count with a revision number. State the watermark delay, state retention and maximum backlog age. Merchant category must match the version effective at payment time, not the category visible when the replay runs.
Build recoverable state
Key totals by account and event hour. Persist source position, deduplication identities, open-window state and reference versions in a consistent checkpoint. Write output candidates under a stable account-hour-revision key, then commit through a sink that rejects repeated keys. Checkpoint design protects the state; sink identity protects reader-visible effects.
Handle disorder and corrections
Emit a provisional window, then send a valid late payment before the allowed cutoff and publish revision two. Send another payment beyond ordinary retention; route it to a signed repair request rather than silently dropping it. Deliver an effective-dated merchant change after a payment and verify the historical join follows the declared temporal rule.
Inject pressure and failure
Throttle the sink until lag reaches a known backlog, restore capacity and measure net drain. Kill a worker after output candidate creation but before checkpoint completion; recover from the previous checkpoint and assert one committed revision per key. Test a duplicate payment ID, overlapping merchant versions and a state serializer mismatch. The last case must stop safely.
Deliver release evidence
Provide source and output schemas, watermark and retention policy, a checkpoint/restore transcript, temporal-join fixtures, revision history, backlog calculation, sink idempotency proof and reconciliation of accepted, quarantined and repaired events. Record the exact source positions and code revision so the release is reproducible.
Implementation
events = [
{"id": "pay-47", "account": "acct-8", "hour": 47, "cents": 2600},
{"id": "pay-48", "account": "acct-8", "hour": 47, "cents": 900},
]
def ledger_snapshot(records):
seen = set()
totals = {}
for row in records:
if row["id"] in seen:
continue
seen.add(row["id"])
key = (row["account"], row["hour"])
totals[key] = totals.get(key, 0) + row["cents"]
return totals
assert ledger_snapshot(events + events[:1]) == {("acct-8", 47): 3500}Performance and operating cost
The reference pass costs O(N) expected time and O(E + K) state for retained identities and active account-hour keys. Production cost also includes checkpoint I/O, temporal-reference versions, broker backlog, output revisions and repair work. Keep a retention limit that covers expected disorder without allowing unbounded state growth.
Common Mistakes
- Do not call a provisional result final.
- Do not replay an external side effect without its idempotency identity.
- Do not equate restored process health with caught-up event-time freshness.
Read next
- Keyed state, checkpoints and recovery
- Watermark retention and late corrections
- Stream-table joins and reference versions
- Stream backpressure and catch-up budget
- Replay manifests and audit trails
Continue the workflow: Event envelopes, identity and tombstones.
Continue the workflow: Project: recover a payment-risk interval join.
Continue the workflow: Project: recover late customer activity sessions.
