The platform had two legitimate but different problems. A high-write PostgreSQL backend produced a large volume of WAL, and business-critical changes took too long to become useful in Snowflake. Treating those as one “CDC scaling” problem would have led to the wrong design.
The original path was durable, not fast
Debezium read PostgreSQL logical replication changes and published them to Pub/Sub. Events were persisted to GCS, then discovered and loaded into Snowflake. The object-storage path was valuable for durability, replay and recovery, but its file boundaries put batching directly on the freshness path.
PostgreSQL
→ WAL / logical replication
→ Debezium
→ Pub/Sub
→ GCS files
→ load discovery
→ Snowflake
→ scheduled transformation
→ usable analytical data
Raw events could be available upstream while useful data still waited for files, ingestion and warehouse transformations. The business requirement was not “copy bytes quickly”; it was to make validated, transformed state queryable in near real time.
Put Spark where continuous processing was needed
Spark Structured Streaming consumed the Pub/Sub event stream continuously, deserialized Debezium envelopes, normalized schemas, handled ordering and duplicates, applied business transformations and wrote micro-batches to Snowflake. GCS remained as the durable replay path, but it left the real-time critical path.
┌─→ GCS
PostgreSQL → Debezium → Pub/Sub archive + replay
│
└─→ Spark Structured Streaming
├─ validate
├─ deduplicate
├─ normalize
├─ enrich
└─ aggregate
↓
Snowflake
Why Spark, rather than only streaming ingestion?
If the requirement had been to insert unchanged events into Snowflake, a native streaming ingestion service would have been simpler. Spark earned its place because the stream needed substantial processing before it became an analytical product.
The critical boundary: Spark does not fix PostgreSQL WAL
Spark was downstream of Pub/Sub. It could not reduce the WAL already produced by PostgreSQL, and it could not make Debezium consume a replication slot faster. If WAL generation exceeds Debezium's sustainable consumption rate, retention grows regardless of downstream Spark throughput.
Downstream speed cannot repair an upstream replication-slot bottleneck.
The upstream protections were separate: monitor slot lag and retained WAL, size Debezium for peak sustainable throughput, filter unnecessary change events, bound retained WAL, and maintain recovery paths that do not depend on preserving days of database logs.
Correctness under at-least-once delivery
A streaming path must expect duplicates, retries and partial failure. Events carried stable identities and source ordering. Stateful deduplication was bounded by event time and watermark policy. Snowflake writes used deterministic keys and idempotent merge behavior. Checkpoints recorded stream progress, but sink idempotency—not a marketing promise of “exactly once”—protected the business result.
event_key = (source_table, primary_key, source_version)
micro_batch
→ reject invalid envelopes
→ keep latest valid source version
→ transform deterministically
→ MERGE on event_key
→ record batch + source offsets
→ advance checkpoint
The outcome
The architecture moved business-critical freshness from twice-daily delivery to near-real-time usable data. It also shifted high-volume transformation away from repeated Snowflake warehouse work. As part of broader platform optimization, Snowflake operating cost fell by 34% while reliability improved.
The important lesson is not “Spark is good for big data.” It is more precise: Spark made sense because a sustained CDC stream needed low-latency transformation, data reduction and controlled delivery before Snowflake. The WAL problem remained a PostgreSQL and Debezium concern—and the design treated it that way.
This case study is intentionally anonymized and simplified. It describes architecture and engineering trade-offs without identifying the company or exposing proprietary data.