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.

Original critical 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.

Separated real-time and replay paths
                              ┌─→ 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.

Latency: transformations ran continuously instead of waiting for the next warehouse task or dbt batch.
Cost control: filtering, deduplication and aggregation reduced the volume Snowflake had to store and repeatedly scan.
Throughput: sustained high-volume processing used explicitly sized distributed compute.
Portability: transformation logic and recovery behavior were not hidden inside one warehouse.

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.

Simplified write contract
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.