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

Project: operate a daily billing pipeline

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

Release a daily billing DAG whose readiness checks, retry rules, lineage, SLO and replay manifest survive a late source and an uncertain export acknowledgement.

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

python
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

Continue the workflow: Project: release a reconciled revenue mart.

Continue the workflow: Project: release a recoverable streaming risk ledger.

ai-data
data-engineering
Storage details