Pular para o conteúdo

apps/data-patterns/change-data-capture

086 · Dados & Persistência · ≈ 6 min de estudo

Change Data Capture (CDC)

Captura mudanças em nível de linha (INSERT, UPDATE, DELETE) de uma tabela e as registra como um fluxo de eventos consumível por sistemas downstream. Aqui o CDC é implementado via trigger: uma função PL/pgSQL grava cada alteração — com estado antigo e novo — na tabela cdc_events. Um publisher real faz a ponte entre esse log de mudanças e um tópico Kafka, de onde consumers de verdade leem o fluxo.

passos
6
arquivos
11
testes
2
tecnologias
5
Infraestrutura realTypeScriptBunPostgreSQLKafkaDocker
Baixar cartão

Cenário

Um core bancário muda contas — limite de crédito, status do cartão — e antifraude, data warehouse e notificações precisam saber, sem que cada um consulte a tabela e sem mexer no código do core. O CDC lê as mudanças do próprio log de transações do PostgreSQL (WAL) e as publica como eventos num tópico Kafka; os consumidores reagem a partir dali.

Planta

Sequência
6/6
INSERT, UPDATE, DELETE em accounts (SQL puro)1pg_logical_slot_peek_changes (só transações commitadas)2evento por linha, key = conta, antes e depois3pg_replication_slot_advance até o último COMMIT publicado4consumer group, commit de offset + 15projeção idempotente pela posição no WAL6Core bancárioPostgreSQL (WAL, wal_level=logical)Relay CDCKafka core.account.changedConsumer do data warehouse
6

6 passos — reproduza para seguir o fluxo

Como funciona

6 passos

Captura mudanças direto do log do banco.

Esconde cada passo: lembre antes de tocar para revelar.

  1. 01

    O PostgreSQL roda com wal_level=logical e a tabela accounts com REPLICA IDENTITY FULL: o WAL guarda a linha inteira de antes em UPDATE e DELETE

  2. 02

    Um slot de replicação lógica (test_decoding) retém o WAL a partir da última posição confirmada — relay fora do ar não perde nada

  3. 03

    O relay lê com peek só transações commitadas (rollback nunca aparece), gera um evento por linha com before, after e changedFields e publica com key = conta e acks: -1

  4. 04

    Só depois do publish confirmado avança o slot até o último COMMIT publicado; se cair no meio, republica — at-least-once

  5. 05

    O consumer group do data warehouse aplica cada evento numa projeção só se a posição no WAL for maior que a gravada — reentrega não muda nada — e commita offset + 1

  6. 06

    Sem trigger e sem tabela de eventos: a escrita do core não paga custo extra

apps/data-patterns/change-data-capture

11 arquivos

sql/

  • 01_schema.sqlschemaTabela de origem com REPLICA IDENTITY FULL e projeção do data warehouse

src/

  • decoding_wal.tsParser da saída do test_decoding: transação por transação, evento por linha
  • relay_cdc.tsSlot, leitura com peek, publish no Kafka e avanço do slot
  • consumer_projection.tsConsumer group com projeção idempotente pela posição no WAL
  • topic_cdc.tsCliente e tópico com partições explícitas
  • pool_cdc.tsPool do PostgreSQL
  • config_cdc.tsVariáveis de ambiente validadas no boot
  • decoding_wal.test.tstesteParser: aspas, nulos, LSN e transação sem commit
  • relay_cdc.test.tstesteIntegração com PostgreSQL e Kafka reais: sem trigger, rollback, republicação após queda e alcance do atraso
  • demo.tsdemoMudanças em SQL puro chegando à cópia do data warehouse

./

  • docker-compose.ymlinfraPostgreSQL 16 com decodificação lógica (5432) e Kafka 3.7 em KRaft (9092)

Executar · com Docker

  1. docker compose up -d --wait# sobe PostgreSQL · Kafka
  2. cp .env.example .env# variáveis de ambiente
  3. bun install# dependências
  4. bun run demo# roda o cenário
  5. bun run test# integração contra o serviço real

Sob Bun o KafkaJS imprime TimeoutNegativeWarning e, com o broker recém-criado, erros de coordinator ainda carregando que ele mesmo repete — os dois são inofensivos.

Requisitos
BunDocker
Sobe junto
PostgreSQLKafka

Por que se relacionam

Teste rápido

Qual destes combina com Change Data Capture?

Próximo projeto · Dados & PersistênciaSnapshot Pattern
Esc

↑ ↓ navegarEnter abrir191 resultados