Release a daily billing DAG whose readiness checks, retry rules, lineage, SLO and replay manifest survive a late source and an uncertain export acknowledgement.
Project: operate a daily billing pipeline
Specify the product
Finance needs one committed bill total per account and event day by 08:30. The grain is account plus day; input orders and exchange rates arrive in separate bounded manifests. A missing source interval blocks publication. An empty but complete order interval is valid and yields a zero-row fact with an explicit completion marker, not an invented carry-forward total.
Build the DAG
Create readiness tasks for both sources, a transform pinned to their snapshot IDs, a uniqueness and sum reconciliation task, a snapshot commit, and an export task keyed by interval and output snapshot. A late rate source should leave yesterday's published snapshot untouched. Idempotent task outputs make a retry converge instead of appending another bill.
Inject three failures
First, send half a source batch without its completion manifest; no transform should run. Second, make one account produce a duplicate business key; quarantine or fail the candidate, keeping the last good snapshot. Third, let the export receiver accept the file but drop the acknowledgement; a retry must use the same idempotency key and produce one remote artifact. Retry policy should stop on the duplicate-key defect.
Evaluate the consumer promise
For each interval, calculate whether complete inputs, quality and publication met 08:30. Keep source lag and processing duration as separate diagnostics. Record a run manifest with input snapshots, code digest, counts, output snapshot and export key; emit lineage to the finance report. The SLO counts a late but correct interval as late, and an on-time wrong interval as bad.
Deliver a release packet
Include an interval diagram, a contract for each source, the retry table, a manifest from one successful run, evidence for all three injected failures, a query that checks one-row-per-account-day, and an operator action for each alert. Run an exact replay from retained inputs and compare the canonical output digest. If a source version is unavailable, name the limit rather than asserting reproducibility.
Implementation
from collections import defaultdict
order_rows = [
{"account": "acct-47", "day": 47, "amount_cents": 7300},
{"account": "acct-47", "day": 47, "amount_cents": 2500},
{"account": "acct-83", "day": 47, "amount_cents": 4100},
]
def bill_totals(records):
totals = defaultdict(int)
for record in records:
if record["amount_cents"] < 0:
raise ValueError("negative order amount requires separate refund flow")
totals[(record["account"], record["day"])] += record["amount_cents"]
return dict(totals)
totals = bill_totals(order_rows)
assert totals[("acct-47", 47)] == 9800
assert sum(totals.values()) == sum(row["amount_cents"] for row in order_rows)Performance and operating cost
Aggregation is O(N) expected time and O(A) space for N orders and A account-day keys. Production cost includes pinned input retention, validation scans, transactional publication, monitoring and occasional replay. A fast transform is not a successful release if its source or export boundary is unverified.
Common Mistakes
- Do not treat an incomplete batch as a legitimate zero-row interval.
- Do not retry an export with a new key after an uncertain acknowledgement.
- Do not call a scheduler-green run good when its published totals fail reconciliation.
Read next
- DAG intervals and idempotent task outputs
- Dependency readiness and source freshness
- Pipeline lineage and impact analysis
- Pipeline SLOs, freshness and error budgets
- Failure classification and retry budgets
- Replay manifests and audit trails
- Project: build a replayable receipt pipeline with an atomic publish gate
Continue the workflow: Project: release a reconciled revenue mart.
Continue the workflow: Project: release a recoverable streaming risk ledger.
