A plataforma tinha dois problemas legítimos e diferentes. Um backend PostgreSQL com muitas gravações produzia grande volume de WAL, e mudanças críticas demoravam para se tornar úteis no Snowflake. Tratar ambos como um único problema de “escala de CDC” levaria ao design errado.

O caminho original era durável, não rápido

Debezium lia mudanças por replicação lógica no PostgreSQL e publicava no Pub/Sub. Os eventos eram persistidos no GCS, descobertos e carregados no Snowflake. O caminho pelo object storage era valioso para durabilidade, replay e recuperação, mas os limites de arquivos colocavam batching diretamente no caminho de atualização.

Caminho crítico original
PostgreSQL
  → WAL / replicação lógica
  → Debezium
  → Pub/Sub
  → arquivos no GCS
  → descoberta e carga
  → Snowflake
  → transformação agendada
  → dados analíticos úteis

Eventos brutos podiam estar disponíveis upstream enquanto os dados úteis ainda aguardavam arquivo, ingestão e transformação no warehouse. O requisito não era “copiar bytes rapidamente”; era tornar estado validado e transformado consultável quase em tempo real.

Colocar Spark onde processamento contínuo era necessário

Spark Structured Streaming consumia continuamente o stream do Pub/Sub, desserializava envelopes do Debezium, normalizava schemas, tratava ordenação e duplicatas, aplicava transformações de negócio e escrevia micro-batches no Snowflake. GCS permanecia como caminho durável de replay, fora do caminho crítico de tempo real.

Caminhos separados de tempo real e replay
                              ┌─→ GCS
PostgreSQL → Debezium → Pub/Sub     arquivo + replay
                              │
                              └─→ Spark Structured Streaming
                                    ├─ validar
                                    ├─ deduplicar
                                    ├─ normalizar
                                    ├─ enriquecer
                                    └─ agregar
                                          ↓
                                      Snowflake

Por que Spark, e não apenas ingestão streaming?

Se o requisito fosse inserir eventos sem alteração no Snowflake, uma ingestão streaming nativa seria mais simples. Spark justificava sua presença porque o stream exigia processamento substancial antes de se tornar produto analítico.

Latência: transformações contínuas substituíram a espera pelo próximo task ou batch do dbt.
Controle de custo: filtro, deduplicação e agregação reduziram o volume armazenado e reprocessado no Snowflake.
Throughput: processamento contínuo de alto volume usou compute distribuído dimensionado explicitamente.
Portabilidade: transformação e recuperação não ficaram escondidas dentro de um único warehouse.

A fronteira crítica: Spark não corrige WAL no PostgreSQL

Spark estava downstream do Pub/Sub. Ele não reduzia o WAL já produzido pelo PostgreSQL e não fazia Debezium consumir o replication slot mais rápido. Se a geração de WAL ultrapassa a capacidade sustentável do Debezium, a retenção cresce independentemente do throughput do Spark.

Velocidade downstream não corrige um gargalo upstream no replication slot.

As proteções upstream eram outras: monitorar lag e WAL retido, dimensionar Debezium para o pico sustentável, filtrar mudanças desnecessárias, limitar a retenção e manter caminhos de recuperação que não dependam de preservar dias de logs do banco.

Correção sob entrega at-least-once

Um caminho streaming precisa esperar duplicatas, retentativas e falha parcial. Eventos carregavam identidades estáveis e ordenação da fonte. Deduplicação stateful era limitada por event time e watermark. Escritas no Snowflake usavam chaves determinísticas e merge idempotente. Checkpoints registravam progresso, mas era a idempotência do sink—não uma promessa abstrata de “exactly once”—que protegia o resultado.

Contrato de escrita simplificado
chave_evento = (tabela_origem, chave_primaria, versao_origem)

micro_batch
  → rejeitar envelopes inválidos
  → manter a versão válida mais recente
  → transformar deterministicamente
  → MERGE pela chave_evento
  → registrar batch + offsets de origem
  → avançar checkpoint

O resultado

A arquitetura levou a atualização de dados críticos de duas vezes ao dia para dados úteis quase em tempo real. Também moveu transformações de alto volume para fora de trabalho repetido no warehouse. Como parte de uma otimização mais ampla, o custo operacional do Snowflake caiu 34% e a confiabilidade melhorou.

A lição não é “Spark serve para big data”. É mais precisa: Spark fazia sentido porque um stream CDC contínuo precisava de transformação com baixa latência, redução de dados e entrega controlada antes do Snowflake. O problema de WAL continuava pertencendo ao PostgreSQL e ao Debezium—e o design respeitava essa fronteira.

Este estudo de caso foi anonimizado e simplificado. Ele descreve arquitetura e trade-offs sem identificar a empresa ou expor dados proprietários.