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

Project: release a recoverable streaming risk ledger

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

Build a payment-risk ledger that joins historical merchant categories, revises late windows and recovers without duplicate published effects.

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

python
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

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.

ai-data
data-engineering
Storage details