PRODPlataforma BaaS europeia soberanaAbra o painel →

Engenharia · 10 minutos de leitura

PostgreSQL CDC com NATS: arquitetura em tempo real

Affane Daylami · Fondateur · 24 de julho de 2026

Voltar ao blog

Um cliente Aurabase conectado via WebSocket vê uma linha inserida na base aparecer alguns milissegundos após o COMMIT — sem polling, sem webhook para configurar. O mecanismo está em três etapas. Um plugin de replicação lógica decodifica o Postgres WAL em JSON, apenas um processo tem o direito de lê-lo e o NATS retransmite cada evento para todos os servidores WebSocket no cluster. Veja como ele está conectado ao código real do aura-realtimeservice.

Este texto em inglês foi gerado automaticamente a partir do original em francês e ainda não foi revisado.
Esta página foi traduzida automaticamente. A versão em inglês é oficial.

Nada aqui é um padrão idealizado. Cada extrato vem da caixa aura-realtime, na data de publicação. Há uma distinção clara aqui: o NATS JetStream não é o que transporta os eventos do CDC neste pipeline. Detalhamos abaixo o que o JetStream realmente faz aqui e o que transporta o tráfego do CDC em seu lugar. Uma armadilha verificada do Kubernetes, capaz de silenciar todo o mecanismo sem gerar erros, está documentada mais abaixo.

O essencial
  • O Postgres CDC da Aurabase lê o WAL via wal2json e pg_logical_slot_get_changes - não o protocolo binário pgoutput, não a ponte Debezium/Kafka Connect.
  • O slot lógico aura_cdc_slot tolera apenas um leitor: aura-realtime-cdc-worker é eleito por meio de um Lease Kubernetes (coordination.k8s.io/v1), com um modo "sempre líder" excluindo K8s.
  • O fan-out para as réplicas ws-front passa pelo núcleo NATS (pub/sub simples em aura.realtime.>), não por meio de um fluxo JetStream persistente. JetStream, neste mesmo serviço, atende exclusivamente KV de presença entre instâncias.
  • Cada evento é verificado novamente em RLS pelo assinante, pouco antes da transmissão, por meio de uma solicitação SET LOCAL ROLE real na linha em questão – não uma aproximação em cache.
  • Armadilha verificada no histórico do repositório: sem o plugin wal2json na imagem Postgres Kubernetes, a criação do slot falha e o CDC permanece silenciosamente inativo.
#
Escolha do motor

wal2json, não pgoutput, não Debezium

A maioria dos pipelines do CDC Postgres passam por pgoutput, o protocolo de replicação lógica binária, e depois por um conector como o Debezium, que o traduz para Kafka. Aurabase pula esta etapa: o serviço aura-realtime decodifica diretamente o WAL em replicação lógica com o plugin wal2json, que produz JSON utilizável sem um estágio de tradução intermediário.

cdc/postgres.rs (consultas reais)sql
-- Criado uma vez se estiver ausente, recriado se o WAL expirar
SELECT pg_create_logical_replication_slot(
  'aura_cdc_slot', 'wal2json');

-- Pesquisa repetida (retirada de 100 ms quando o lote retorna vazio)
SELECT data FROM pg_logical_slot_get_changes(
  'aura_cdc_slot', NULL, 1000,
  'include-xids', '0', 'format-version', '2',
  'include-schemas', '1', 'include-types', '0');

Três pré-requisitos do Postgres, verificados no código: a função do aplicativo deve carregar REPLICATION, wal_level = logical deve estar ativa e max_slot_wal_keep_size deve limitar a retenção do WAL. Sem este último limite, um consumidor lento faz com que o volume do disco cresça indefinidamente. O serviço também monitora wal_status do slot: se ele mudar para lost (WAL eliminado além deste limite), o slot será recriado automaticamente. A perda de eventos não consumidos entretanto é assumida e registrada.

Um lote de 1.000 alterações por enquete, com um filtro de tabela dinâmico atualizado a cada 3 segundos. Somente tabelas onde um projeto habilitou explicitamente o tempo real entram em add-tables — isso evita a decodificação do WAL de tabelas que não são de interesse de nenhum assinante.

#
Um único produtor

cdc-worker: eleito por um arrendamento Kubernetes, não duplicado

Um slot de replicação lógica Postgres tolera apenas uma unidade ativa por vez. Executar vários consumidores do mesmo aura_cdc_slot em paralelo quebraria a ordem das alterações, e não apenas as duplicaria. Aurabase dividiu aura-realtime em dois binários separados para resolver essa restrição sem sacrificar a escala horizontal do WebSocket.

Cargo.tomltoml
# cdc-worker: líder eleito, lança o CDC + publica no NATS. Nenhum servidor WS.
[[bin]]
name = "aura-realtime-cdc-worker"

# ws-front: N réplicas (HPA), consome NATS, filtro RLS, atende WS/SSE.
[[bin]]
name = "aura-realtime-ws-front"

cdc-worker só ganha o direito de ler o slot mantendo um objeto Lease Kubernetes (coordination.k8s.io/v1). Ele tenta criá-lo ou roubá-lo se ele tiver expirado e então o renova após um terço de sua vida útil. Quando o arrendamento é perdido (falha na renovação, mudança de titular), um CancellationToken corta o trabalho em andamento e o processo retorna a um ciclo de aquisição. Fora do Kubernetes – durante testes locais ou desenvolvimento sem cluster, por exemplo – o cliente K8s não consegue se conectar e o serviço muda para o modo “sempre líder”. Útil em desenvolvimento, incorreto se você esquecer em produção com várias réplicas.

cdc/leader.rs (assinaturas reais)rust
pub async fn run_with_lease<F, Fut>(cfg: LeaseConfig, work: F) -> anyhow::Result<()>
where F: Fn(CancellationToken) -> Fut, ...
{
  // renova todos oslease_duration_secs/3
  // excluindo K8s: funciona(token) apenas uma vez, token nunca cancelado
}
#
Fan-out real

Núcleo NATS para CDC, JetStream para presença

Este é o ponto mais mal compreendido neste tipo de pipeline: NATS e JetStream são duas coisas distintas, e os eventos CDC não passam por um fluxo JetStream persistente. cdc-worker publica cada evento com Client::publish_with_headers — a API principal pub/sub do NATS, aquela entregue no máximo uma vez sem persistência ou repetição, separada do JetStream.

Coração NATS (pub/sub)Distribuição de eventos do CDC para réplicas ws-frontNo máximo uma vez, sem persistência ou repetição
NATS JetStream (KV)Presença entre instâncias (aura_presence)Estado compartilhado dos anos 60, não um fluxo do CDC para reproduzir
channels/manager.rs → nats/mod.rs (extraits réels)rust
// cdc-worker publica (ChannelManager::publish_and_relay)
nc.publish_with_headers(subject, headers, payload).await

// ws-front escuta toda a subárvore (NatsRelay::start_relay)
client.subscribe("aura.realtime.>").await

Cada réplica ws-front se inscreve na subárvore aura.realtime.>e decodifica o tópico para o nome do canal original. Em seguida, ele republica o evento em seu tokio::sync::broadcastlocal, para os WebSockets conectados a ele. Um cabeçalho Aura-Origin carrega o identificador da instância remetente: cada réplica ignora as mensagens que ela mesma publicou, o que evita duplicatas sem coordenação central.

Esta escolha tem uma consequência direta: o núcleo NATS entrega no máximo uma vez (at-most-once), sem uma fila de longa duração. Se uma réplica ws-front for brevemente desconectada do NATS quando um evento passar, ela não será atualizada. O histórico por canal (Channel.history, um buffer na memória limitado por HISTORY_SIZE) reside localmente em cada réplica, não no nível do cluster. Este é um compromisso aceitável precisamente porque o slot de replicação lógica do Postgres continua sendo a fonte de verdade durável upstream. É wal2json que garante que nenhuma alteração no banco de dados seja perdida antes do consumo, não no NATS.

JetStream existe em aura-realtime — mas para um uso totalmente diferente. Ele atende o armazenamento KV de presença entre instâncias (aura_presence, armazenamento de memória, max_agede 60 segundos ), que sincroniza quem está online em qual canal entre todas as réplicas ws-front. Um estado compartilhado de curta duração, não um fluxo de eventos do CDC para reproduzir.

#
Segurança

RLS verificado novamente pelo assinante, em cada evento

Um evento do CDC não é transmitido para todos os assinantes de um canal. O poller roteia cada alteração para um dos trabalhadores CDC_NUM_SHARDS (4 por padrão, com hash de project_id), que primeiro pesquisa pg_policies. Uma tabela sem uma política RLS permite todos os assinantes — o comportamento Supabase — enquanto uma tabela com políticas aciona uma verificação por assinante.

subscriptions.rs::check_rls_visible_batch (real, simplificado)sql
BEGIN;
SET LOCAL ROLE "aura_authenticated";
-- Reivindicações JWT do assinante injetadas no GUC
SELECT set_config('request.jwt.claims', $1, true);
SELECT EXISTS(SELECT 1 FROM "schema"."table" WHERE id::text = $2);
ROLLBACK;

A transação assume a função Postgres física do assinante e injeta suas declarações JWT em request.jwt.claims — aquele que auth.uid() e auth.role() leem no lado da política. Em seguida, verifica se a linha permanece visível nesta função e cancela: sem gravação, uma leitura RLS real. Somente sub_id validados pousam em allowed_sub_ids do evento; um vec vazio significa ninguém, em fechado com falha.

Astuce

Um DELETE escapa dessa verificação por linha - a linha desapareceu, é impossível relê-la em RLS - e volta para todos os assinantes do canal. Em troca, a carga útil de um DELETE expõe apenas as colunas de identidade (chave primária), nunca o conteúdo da linha excluída.

Um detalhe que muitas vezes atrapalha na migração: wal2json só captura o conteúdo completo de um UPDATE ou de um DELETE se a tabela tiver REPLICA IDENTITY FULL. Sem ele, apenas o INSERT volta a funcionar. É por isso que o endpoint que permite tempo real em uma tabela define este ALTER TABLE automaticamente — uma correção real para o repositório, não uma caixa de seleção separada.

#
Armadilha verificada

Por que o CDC pode permanecer silencioso no Kubernetes

A história do processo documenta uma armadilha real, não um caso teórico. A imagem Postgres usada por padrão em um cluster Kubernetes (a imagem da comunidade de estoque pgvector/pgvector) não inclui o plug-in wal2json. Sem ele, pg_create_logical_replication_slot('aura_cdc_slot', 'wal2json') falha e nada no lado do cliente sinaliza isso explicitamente. Os WebSockets permanecem abertos, as assinaturas são aceitas, mas nenhum evento de banco de dados acontece.

docker/Postgres.Dockerfile (extrato real)dockerfile
RUN apt-get update && \
    apt-get install -y --no-install-recommends \
      postgresql-16-wal2json=2.6-4.pgdg12+1 \
      ...

A correção real foi construir e publicar uma imagem Postgres dedicada com este pacote instalado e, em seguida, apontar o Kubernetes StatefulSet para ela em vez da imagem de estoque. Ele também definiu explicitamente wal_level=logical, max_replication_slots e max_slot_wal_keep_size como argumentos de inicialização - ausentes em imagens genéricas.

É precisamente para tornar esse tipo de falha observável que cdc-worker expõe um sinalizador cdc_active (falso, desde que o slot não seja confirmado como íntegro) em seu endpoint /health dedicado. Este é um sinal binário a ser observado – em vez de inferir CDC morto apenas pela ausência de eventos do lado do cliente.

#
Limite atual

O que este pipeline ainda não cobre: clusters dedicados

Este mecanismo lê uma única variável POSTGRES_REPLICATION_URL, portanto, um único servidor Postgres e um único slot lógico. Isso é consistente com o esquema multilocatário por modelo de projeto (project_<uuid>) em um cluster compartilhado. Nesse modelo tudo funciona: um único cdc-worker vê as alterações de todos os projetos do mesmo cluster e roteia por diagrama.

Para um projeto em um cluster CNPG dedicado - sua própria instância isolada do Postgres - a imagem usada hoje (construída a partir da imagem de estoque CloudNativePG) adiciona apenas pg_graphql, não wal2json. Nada também executa um cdc-worker por cluster dedicado. No estado atual do código, o CDC em tempo real descrito aqui permanece, portanto, focado na camada compartilhada – um limite arquitetônico real, não uma funcionalidade de roteiro.

#
Lado do cliente

O que não muda: a API postgres_changes do SDK

Todo esse mecanismo permanece invisível no SDK. A API encadeadaaurabase-js não muda dependendo se o evento vem de cdc-worker via NATS ou, no desenvolvimento local, do binário aura-realtime não dividido.

app.tstypescript
const channel = aura.realtime.channel('room-1')
  .on('postgres_changes', {
    event: 'INSERT', schema: 'public', table: 'messages'
  }, (payload) => { /* … */ })
  .subscribe()

Mas antes que um único evento possa acontecer, a tabela deve ser registrada no CDC – isso não é automático na criação. Uma chamada service_role no endpoint dedicado (ou a alternância equivalente no Studio) grava a tabela em realtime_tables do esquema da plataforma do projeto e define REPLICA IDENTITY FULL para você.

terminalbash
curl -X PUT \
  https://api.aurabase.cloud/v1/db/<project_id>/schema/tables/messages/realtime \
  -H "apikey: <service_role_key>" \
  -H 'Content-Type: application/json' \
  -d '{"enabled": true}'

Para o restante da superfície SDK compatível com Supabase — autenticação, armazenamento, políticas RLS —, consulte o guia de migração. Para saber como aura-realtime se divide em vários binários em uma única caixa de espaço de trabalho, consulte o detalhe da arquitetura de carga .

PRONTO PARA IMPLEMENTAR?

Seu back-end em cinco minutos.

Não é necessário cartão de crédito · 500 MB grátis · 50.000 MAU