Finding and Fixing Silent Data Loss in a Production Lead Pipeline
A production lead-generation pipeline scraping ~6,400 businesses twice daily, opened with a formal data-quality audit. Found two critical, silent defects — a dead dedup key that destroyed 545 distinct businesses, and a token-budget bug that erased 57.8% of one scraper's output — then rebuilt the pipeline as a Fabric-shaped lakehouse with 197+ automated assertions so the same class of bug can't recur silently.
Personal, self-directed · Data Engineering · Published 2026-09-07
Results
- ▲545 distinct businesses recovered from a dead dedup key silently collapsing 8.5% of the table
- ▲57.8% of one scraper's enrichment output recovered from a token-budget bug indistinguishable from a legitimate empty result
- ▲A remediation workflow caught in review before it ran — it would have inserted up to 6,378 duplicate rows per execution
- ▲A placeholder join caught inflating a fact table 28x, invisibly, through a fully green dbt build
- ▲197+ automated assertions now gate every layer (unit / compiled-workflow / cross-workflow-parity / live-model integration)
- ▲Lakehouse runs end to end on a live, armed Dagster schedule — not scaffolded — with dbt build clean at 36 pass / 0 warn / 0 error
I run a production lead-generation pipeline that scrapes ~6,400 Philippine businesses twice daily, scores them against a loan-company ICP with a local LLM, and hands ~2,000 qualified leads to a sales team. I'm a software engineer, not a trained data engineer — I built this because a real pipeline needed one, using tools I already knew (n8n, Postgres, a local LLM), not a resume-shaped tech stack. The medallion-architecture vocabulary I picked up *after* building something that turned out to already have that shape.
The Audit — Where the Real Engineering Is
I opened a full rebuild with a formal data-quality audit against the live 6,378-row / 2,081-row population — reproducible end-to-end, no estimated figures. It found 13 findings. Two were critical, and **both were silent**: neither raised an error, failed an execution, or fired an alert.
The Dead Key That Ate 545 Businesses
The dedup logic keyed on `place_id`, falling back to `domain`. The scraper never emitted `place_id` — not sometimes empty, absent from the response object entirely. Every run silently used the domain fallback, encoding the assumption "two businesses sharing a website are the same business." For Philippine SMEs whose Google Maps listing points at a shared Facebook page, that assumption was badly wrong: 231 domains were shared by more than one business, and replaying the corrected rule against all 6,378 rows showed **545 distinct businesses had been wrongly destroyed** — 8.5% of the table, silently, on every run, for as long as the bug existed.
The Token Budget That Erased 57.8% of One Scraper's Output
Two scraper workflows shared enrichment logic but diverged on one config value: `max_tokens: 300` vs. `40000`. The model's real responses run 772–775 tokens. At 300, more than half of every response was truncated mid-JSON, caught by a generic parse-error handler, and stored as "no categories apply" — byte-for-byte indistinguishable from a legitimate empty result. Over half of one scraper's entire output was never actually categorized, and nothing reported it, because a parse failure and a genuine zero looked identical downstream.
The Remediation Workflow That Would Have Made It Worse
Fixing the token budget didn't repair the rows already damaged, so I built a remediation workflow — cloned from the original scraper and rewired to read stale rows back out instead of scraping fresh ones. The read side was rewired correctly. The write side wasn't: it still carried the scraper's `create` operation against a table whose only constraint is a UUID primary key. Running it once would have inserted up to 6,378 duplicate rows; on its 4-hour schedule, indefinitely many — and nothing would have caught it, because there was no unique constraint for the error handler (present to absorb a *different* class of error in the original) to absorb.
I caught this in review, before the workflow ever ran, by asking whether a guard that made sense in the original context still did anything once the code was copied into a new one. Fixed: an idempotent update filtered on row id, with three terminal statuses so a re-run is safe. 33 regression assertions pin this specific bug shape — a guard that looks load-bearing but silently does nothing once copied — since it recurs throughout the audit.
The Cross Join That Inflated a Fact Table 28x, Invisibly
A placeholder join in the dbt star schema (`on 1 = 1`, commented "location extracted from address during ingestion") was a cross join against a 28-row dimension table. Every scoring event was emitted 28 times — 52,696 rows for 1,882 real events — and every downstream count, rate, and average was wrong by exactly 28x. `dbt build` reported clean because the marts had **no schema tests at all**: a fan-out doesn't violate a `not_null` or `accepted_values` check, only a uniqueness assertion on the actual grain would have caught it.
What Was Built to Make This Class of Bug Structurally Harder to Reintroduce
- **A written data contract** — one authoritative rule set instead of five unlinked definitions of "what counts as a hot lead" scattered across a prompt, a parser, two pipeline scripts, and a dbt file. Generated code and docs derive from one rules block, checked in CI.
- **A quarantine-rate gate, not a silent drop** — records failing the contract go to a quarantine table with a reason, never discarded and never coerced, gated by a Dagster asset check pinned against the historical 11% quarantine rate.
- **A 4-tier, 197+ assertion test suite** (unit / compiled-workflow / cross-workflow-parity / live-model integration) — n8n stores workflow logic as an escaped JSON string with no linting or type-checking by default, so the project extracts it to ordinary `.js` files, tests those, and injects them at build time.
- **A CI-gated deployment tool** for the workflow engine itself, after finding the existing sync tooling's two documented safety claims (auto-backup, post-push verification) had never actually been implemented.
- **A live, armed Dagster schedule** — found, the hard way, that flipping a config flag to RUNNING describes what a daemon does once one exists to read it, and doesn't itself produce a daemon. Fixed with an explicit workspace, a committed instance config, and a liveness-checked startup script.
Closing the Loop, September 2026
Two findings stayed open by choice after the initial rebuild: `place_id` was never captured upstream (the root cause behind the dead dedup key), and 89 duplicate businesses were already resident in the table from repeated scraper runs. Both have since moved. The scraper now parses Google's stable place identifier straight out of the listing's own link — no workflow change needed, since the downstream pipeline had already been reading and writing that field for months, just never receiving a value. The 89-row reconciliation is built as a transaction-safe merge that re-points scoring history before deleting a duplicate (so a cascade delete never takes a lead's score history down with it), unit-tested independent of any database connection.
Trying to run that reconciliation against the live database surfaced a third thing worth naming honestly: the production database stopped resolving entirely mid-fix — a distinct failure from anything the original audit went looking for, and a gap in its own right, since nothing in the pipeline's freshness checks distinguishes "the source is stale" from "the source is gone." Recorded and backlogged rather than smoothed over, in keeping with the rest of this project.
What This Demonstrates
- **Data contract design and schema evolution** — one authoritative rule set instead of five drifting copies.
- **Ownership of architecture, not just implementation** — every fix traces back to a design decision, not a typo.
- **Finding bugs before they cost anything** — the remediation-workflow defect was caught in review, before a single execution.
- **Testing discipline applied to a low-code tool that actively resists it** — n8n's JSON-embedded JavaScript has no linting or type-checking by default; the project built the missing scaffolding.
- **Honest failure modes** — every "success" in this system (a green dbt build, a passing test suite, a RUNNING schedule flag) was, at some point, found to be reporting success for work it hadn't actually done.
This is personal, self-directed work — not a client or employer deliverable, not sold or licensed. It's built on my own pre-existing lead-pipeline data and deliberately shaped to mirror a Fabric-style medallion architecture because that's the stack my day job uses; no employer code, data, or confidential material was used in building it.