Implement a receipt pipeline that survives retries, deletes, late events and a failed backfill without exposing inconsistent aggregates.
Project: build a replayable receipt pipeline with an atomic publish gate
Define the deliverable
Build a raw landing table, current receipt table and daily aggregate. Each input needs source identity, event time, arrival time and source position. Write correction and delete semantics before building the transform. A second run over identical input must leave exactly the same current state and aggregate. Begin with the source contract], then document retention and ownership.
Implement replay
Consume a deliberately shuffled fixture through a source-order field. Apply an update and later delete for the same receipt; replay of an older update must not resurrect it. Commit the checkpoint with target state or document an equivalent recovery design. CDC ordering] and idempotent writes] define the failure cases.
Test time and publication
Insert a receipt after its business day closes and state whether the daily aggregate is corrected, versioned or left final. Stage a backfill with one required partition missing. The visible aggregate must remain the prior complete version; the failed candidate and reason must remain inspectable. The publish boundary] is the central acceptance test.
Submit evidence
Provide schema, fixture, runnable command and outputs from initial load, same-input replay, crash-and-restart, delete, late arrival and failed backfill. Reconcile source, accepted, rejected and skipped counts. State per-page cost and temporary backfill storage. A final dashboard screenshot cannot establish these invariants.
Implementation
def assert_replay_invariants(load_fixture, snapshot_state, fixture):
load_fixture(fixture)
first_state = snapshot_state()
load_fixture(fixture)
second_state = snapshot_state()
assert first_state == second_state
assert second_state["source_count"] == (
second_state["accepted_count"] + second_state["rejected_count"]
+ second_state["skipped_count"]
)Performance and operating cost
One pass over N source events is O(N) before index costs; a keyed B-tree upsert adds about O(log K) per event for K current receipts. Rebuilding H historical rows costs O(H) plus staging storage.
Common Mistakes
- Do not call a pipeline replayable because its second run finished.
- Do not lose deletes when deriving current state.
- Do not expose a failed staged backfill.
Read next
- Data source contracts: preserve raw records before transformation
- Change data capture: apply updates and deletes by source order
- Event time and late arrivals: close windows with an explicit correction policy
- Backfills: rebuild history without exposing a half-written result
Continue the workflow: Project: release a versioned order lake.
