Build an account-day invoice rollup that handles skew, bounded lookups, task retries and one atomic publication boundary.
Project: release a distributed invoice rollup
Define the release unit
A finance consumer reads one total per account and event day from an immutable source snapshot. Input invoice lines carry invoice ID, account, merchant ID, day and integer cents. The merchant lookup has one row per key and a measured broadcast size. Write the target grain and uniqueness assertion before selecting a compute framework.
Build the execution plan
Filter invalid lines into quarantine, enrich valid lines with the bounded merchant lookup, locally combine cents, and salt only the observed heavy account. Merge salts back to account-day. Capture mapper output count, shuffle bytes, largest reducer bytes and unmatched merchant count. Combining and salting solve different parts of the plan.
Publish once
Write candidate partition files under run and attempt IDs. Validate one row per target key, input-to-output total cents and completeness of every expected partition. A coordinator publishes one manifest referencing accepted files. A reader pinned to the previous manifest must never see half the new rollup. Attempt selection prevents duplicate output after a lost acknowledgement.
Inject failures
Duplicate a merchant lookup key and confirm the run stops before enrichment. Kill one task after its candidate file appears, then retry and verify only one accepted file. Make the heavy account hold 61% of rows and compare p95 task time with and without salting; preserve identical totals and target keys.
Deliver evidence
Include the source and lookup contracts, a physical plan before and after tuning, measured task distributions, one published manifest, a reconciliation report and an exact retry from pinned inputs. If the production engine changes its plan through adaptive execution, record the executed plan rather than only the submitted SQL.
Implementation
from collections import defaultdict
invoice_lines = [
{"id": "inv-47", "account": "acct-47", "day": 47, "cents": 3200},
{"id": "inv-48", "account": "acct-47", "day": 47, "cents": 4700},
{"id": "inv-83", "account": "acct-83", "day": 47, "cents": 2100},
]
def release_totals(lines):
if len({line["id"] for line in lines}) != len(lines):
raise ValueError("duplicate invoice identity")
totals = defaultdict(int)
for line in lines:
totals[(line["account"], line["day"])] += line["cents"]
if sum(totals.values()) != sum(line["cents"] for line in lines):
raise ValueError("reconciliation failed")
return dict(totals)
assert release_totals(invoice_lines)[("acct-47", 47)] == 7900Performance and operating cost
The local reference calculation is O(N) expected time and O(A) memory for N lines and A account-day keys. A distributed release also pays for shuffle bytes, attempt files, validation scans and manifest storage. Tune to the slowest task and preserve the same published result.
Common Mistakes
- Do not optimize away duplicate-key checks to make a join faster.
- Do not compare plans from different input snapshots as if only code changed.
- Do not let task attempts write directly into the reader-visible partition.
