The worst pipeline bug I've shipped didn't throw. No stack trace, no failed job, no alert. Every run went green, finished inside its window, and wrote roughly one record out of every thousand it was handed.
The job's own metrics said it was healthy, because the job's own metrics counted runs, not rows. It took a business user asking why a campaign list looked short to surface it.
This post is about that class of bug — the pipeline that succeeds at doing nothing — and the three habits that make it structurally impossible rather than merely unlikely.
How you lose 99.9% of your data without noticing
Simplified, the stage looked like this. Records arrive in batches, get enriched against a reference lookup, and are staged before reconciliation downstream.
// the bug, reduced
String cacheKey = record.getAccountId();
Enrichment e = cache.get(cacheKey);
if (e == null) {
e = lookupService.fetch(record);
cache.put(cacheKey, e);
}
stage(record, e);
Spot it? accountId isn't unique per record. One account generates many records per batch — different products, different billing cycles, different dates. The key collapsed thousands of distinct records onto a handful of cache entries, and the downstream staging write was keyed off the same value. Last write wins. Everything else evaporated.
The job processed every record. It just persisted almost none of them.
The fix is one line, and the one line is not the lesson:
String cacheKey = String.join("|",
record.getAccountId(),
record.getProductCode(),
record.getBillingCycle(),
record.getRecordDate().toString());
The lesson is that nothing in the system was capable of noticing. That's the actual defect.
Habit 1: count rows at every boundary, and assert on the ratio
A pipeline stage should know how many records it received and how many it emitted, and should refuse to call itself successful when those numbers disagree in a way nobody authorised.
@dataclass
class StageResult:
stage: str
received: int
emitted: int
rejected: int
reason_counts: dict[str, int]
def assert_sane(self, min_ratio: float = 0.95) -> None:
if self.received == 0:
return
accounted = self.emitted + self.rejected
if accounted != self.received:
raise PipelineError(
f"{self.stage}: {self.received - accounted} records vanished"
)
if self.emitted / self.received < min_ratio:
raise PipelineError(
f"{self.stage}: emitted {self.emitted}/{self.received} "
f"({self.emitted / self.received:.1%}) — below floor. "
f"rejections: {self.reason_counts}"
)
Two distinct checks doing different jobs. The first is conservation: every record either came out or was explicitly rejected with a reason. A record that is neither emitted nor rejected has vanished, and vanishing is always a bug. The second is a ratio floor: even when everything is accounted for, a sudden collapse in throughput means something upstream changed.
Deliberate filtering goes through rejected with a reason. That way "we dropped 40% because they were test accounts" is visible and intentional, and "we dropped 40% because of a cache key" is a hard failure.
I now treat a stage without row-conservation accounting as untested, regardless of how many unit tests it has. My own cache-key bug passed its unit tests, because the test fixture had one record per account.
Habit 2: make reruns boring
The question that separates a pipeline you can operate from one you can't: what happens if I run yesterday's batch again, right now?
If the answer involves checking anything first, the pipeline is fragile, because reruns are not an edge case. Reruns happen after an outage, after a bad upstream file, after a bug fix, and always under time pressure with someone asking when it'll be done.
Idempotency isn't a property you bolt on. It comes from one decision: every record has a deterministic natural key, and writes are keyed on it.
INSERT INTO staging_campaign_records (
record_key, account_id, product_code, billing_cycle,
record_date, amount, loaded_at, batch_id
)
VALUES (:record_key, :account_id, :product_code, :billing_cycle,
:record_date, :amount, SYSTIMESTAMP, :batch_id)
ON CONFLICT (record_key) DO UPDATE SET
amount = EXCLUDED.amount,
loaded_at = EXCLUDED.loaded_at,
batch_id = EXCLUDED.batch_id;
The record key is derived from the business meaning of the row — the combination of fields that genuinely identifies it — never from a sequence, a load timestamp, or a UUID generated at read time. Those are all "different every run", which is the opposite of what you need.
Then, crucially, a unique constraint on record_key in the database. Not a check in application code. Application checks race; a unique index is a promise the storage engine keeps even when two workers process overlapping batches.
Note what this gives you beyond safety: deleting and reloading a day's data becomes a thing you can do at 2 p.m. on a Tuesday without ceremony. That changes how fast you can fix things.
Habit 3: let history be history (SCD Type 2 without tears)
The other silent corruption is overwriting a dimension in place. Customer changes segment; you UPDATE the row; every report that ever referenced the old segment is now retroactively wrong, including the ones already circulated.
Type 2 slowly changing dimensions solve this by never updating a fact, only closing it:
-- close the current row when the tracked attributes actually changed
UPDATE dim_customer
SET valid_to = :effective_from - INTERVAL '1 day',
is_current = FALSE
WHERE customer_id = :customer_id
AND is_current = TRUE
AND (segment, tariff_plan, status)
IS DISTINCT FROM (:segment, :tariff_plan, :status);
-- open a new one
INSERT INTO dim_customer (
customer_id, segment, tariff_plan, status,
valid_from, valid_to, is_current
)
SELECT :customer_id, :segment, :tariff_plan, :status,
:effective_from, DATE '9999-12-31', TRUE
WHERE NOT EXISTS (
SELECT 1 FROM dim_customer
WHERE customer_id = :customer_id
AND is_current = TRUE
);
IS DISTINCT FROM is doing real work there — it's null-safe, so a genuine NULL → 'value' transition is detected rather than silently skipped the way != would skip it. And the WHERE NOT EXISTS on the insert makes the pair idempotent: run it twice with the same input and you get one version, not two.
Then querying as-of any point in time is trivial, which is the whole point:
SELECT f.*, d.segment
FROM fact_usage f
JOIN dim_customer d
ON d.customer_id = f.customer_id
AND f.usage_date BETWEEN d.valid_from AND d.valid_to;
Reports stop changing their answers about the past. Disputes become checkable.
The check that would have caught me
Row conservation, as a single assertion, at every stage boundary. That's it. The cache-key bug would have failed its very first run with emitted 1,204/1,100,000 (0.1%) — below floor, instead of looking healthy for weeks.
Everything else in this post is about making that assertion safe to act on: if reruns are idempotent, failing a batch loudly costs you nothing, so you can afford to fail loudly on anything suspicious. A pipeline that can't be safely rerun is a pipeline under pressure to keep going when it shouldn't — and that pressure is where silent data loss actually comes from.
- Count what goes in and what comes out. A record that is neither emitted nor rejected has vanished.
- Make reruns boring. Deterministic natural keys, upserts, unique constraints in the database.
- Never overwrite history. Close the old row, open a new one.
None of it is clever. That's rather the point — clever is what ate my records.
More on backend systems, pipelines and reliability as I go. If you've got a favourite silent-failure story, I'd genuinely like to read it.
Top comments (0)