Back to blog
AutomationSeptember 1, 2026

We Survived a Data Quality Crisis: 4 Hours to Fix a Broken Revenue Pipeline

Survived data quality crisis; fixed broken revenue pipeline in 4 hours

We Survived a Data Quality Crisis: 4 Hours to Fix a Broken Revenue Pipeline

An engineering team noticed it first. Their event pipeline — a simple HTTP endpoint that accepted JSON events and wrote them to a PostgreSQL database — was dropping events during traffic spikes. Black Friday. Product launches. Email blasts. Anytime traffic went above [N] events/second, the endpoint started returning 502s and events disappeared.

The pipeline was built [N] years ago when the company was doing [N] events/day. Now they were doing [N] events/day. The architecture hadn't changed. The database was the bottleneck.

The Diagnosis

The architecture was a single Node.js Express server writing synchronously to a shared PostgreSQL instance:


User Browser

    |

    | POST /events

    v

Node.js Express Server (single instance)

    |

    | INSERT INTO events

    v

PostgreSQL (single instance, shared with app)

Three problems: synchronous writes (the server waited for PostgreSQL to confirm the INSERT before returning 200), single instance (one Node.js process on one server, no horizontal scaling), and shared database (PostgreSQL was also serving the application, so event writes competed with user queries for disk I/O).

The numbers were brutal:

MetricNormalPeakIssue
Events/second[N][N]10x spike
Response time (p99)[N]ms[N]msTimeout
Error rate[PCT]%[PCT]%Events lost
Database CPU[PCT]%[PCT]%Saturated
Events queued0[N]Dropped on restart

The queue wasn't persistent. When the Node.js process restarted (deployments, crashes), queued events were lost. On Black Friday, they lost an estimated [N] events — including [N] purchase events.

What We Built

We rebuilt the pipeline with streaming and decoupling:


User Browser

    |

    | POST /events (async, 200ms timeout)

    v

Load Balancer (ALB)

    |

    v

Event Collector (3x Node.js, auto-scaled)

    |

    | Publish to Kafka (async, 50ms)

    v

Kafka Cluster (3 brokers, replicated)

    |

    | Consume and process

    v

Stream Processors (Flink, 4 workers)

    |

    | Enrich, validate, transform

    v

PostgreSQL Warehouse (dedicated, read-optimized)

    |

    | Query

    v

Analytics (Metabase, dbt models)

Key decisions:

The collector returns 200 immediately after validating the event, before writing to Kafka. The client doesn't wait for persistence. Response time dropped from [N]ms to [N]ms.

Kafka persists events for [N] days. If the downstream processors are slow or down, events accumulate in Kafka. Nothing is lost.

We used Apache Flink for real-time enrichment: user profile lookup (from Redis cache), sessionization (assign session IDs based on 30-min gaps), geo-IP lookup, device parsing, and bot detection.

The PostgreSQL warehouse is separate from the application database. Optimized for analytics: partitioned by date, columnar storage for large tables (using TimescaleDB), read replicas for analytics queries, no application traffic competing for I/O.

The Migration: Zero-Downtime Cutover

We couldn't afford downtime. Their events are business-critical.

Phase 1: Parallel Pipeline (Week 1)

Built the new pipeline alongside the old one. Events went to both. The old pipeline stayed synchronous. The new pipeline was async and validated.

Phase 2: Validation (Weeks 2-3)

Compared old and new pipelines: event counts (within [PCT]%), event contents ([PCT]% match), latency (new: [N]s, old: [N]s), error rates (new: [PCT]%, old: [PCT]%), and data freshness (new: [N]min, old: [N]min).

Found and fixed [N] discrepancies: timezone handling (old used server time, new used UTC), null vs. empty string (old stored "", new stored null), and array serialization (old used comma-separated, new used JSON).

Phase 3: Cutover (Week 4)

Switched primary pipeline to new. Kept old as backup for [N] days.

Phase 4: Cleanup (Week 5)

Removed old pipeline after [N] days of validation.

The Results

MetricBefore (Sync PostgreSQL)After (Kafka + Flink)Change
Events/second (peak)[N][N]+[PCT]%
Response time (p99)[N]ms[N]ms-[PCT]%
Error rate (peak)[PCT]%[PCT]%-[PCT]%
Events lost/day[N]0Eliminated
Data freshness[N]min[N]min-[PCT]%
Infrastructure cost$[N]/mo$[N]/mo+[PCT]%
Engineering time on pipeline[HOURS]/week[HOURS]/week-[PCT]%

The cost increase was [PCT]%, but the business value was: no more lost purchase events (worth $[REVENUE]/year in attribution accuracy), engineering team freed up for product work ([HOURS] hours/week), and Black Friday handled without incident (previously a [N]-hour firefight).

What We Learned

Async is scary until you measure the alternative. The old pipeline felt "safe" because it confirmed writes. But it was losing [N] events/day under load. The async pipeline loses zero events because Kafka persists everything.

Kafka is not a database. It's a buffer. We don't query Kafka directly. Events live in Kafka for [N] days, then are consumed into the warehouse. The warehouse is the source of truth for analytics.

Flink is overkill for simple enrichment. If we were doing this again, we'd start with Kafka Connect + simple consumers. Flink is great for complex stream processing (windowing, joins, stateful operations). For simple enrichment, it's unnecessary complexity.

Validation is a separate concern. We run validation in the collector (fast, reject bad events) and in the stream processor (thorough, route invalid events to dead letter queue). The collector validation is for client feedback. The processor validation is for data quality.

Monitoring is not optional. We track Kafka lag (how far behind consumers are), consumer group health (are all partitions being consumed?), event latency (time from browser to warehouse), dead letter queue size (how many invalid events), and schema compliance rate (how many events match expected schema).

Bottom Line

This client's event pipeline went from a single-point-of-failure to a resilient streaming architecture. The migration took [N] weeks, cost $[N] in infrastructure, and eliminated the event loss that was costing them $[REVENUE] in attribution accuracy.

If your event pipeline is a single HTTP endpoint writing to a database, you have this problem. The fix is streaming, and it's not as complex as it sounds.