PRODPiattaforma BaaS europea sovranaApri Dashboard →

Ingegneria · 10 lettura minima

PostgreSQL CDC con NATS: architettura in tempo reale

Affane Daylami · Fondateur · 24 luglio 2026

Torniamo al blog

Un client Aurabase connesso tramite WebSocket vede apparire una riga inserita nella base pochi millisecondi dopo il COMMIT — senza polling, senza webhook da configurare. Il meccanismo è in tre fasi. Un plug-in di replica logica decodifica Postgres WAL in JSON, solo un processo ha il diritto di leggerlo e NATS inoltra ogni evento a tutti i server WebSocket nel cluster. Ecco come è collegato nel codice attuale di aura-realtimeservice.

Questo testo inglese è stato generato automaticamente dall'originale francese e non è stato ancora rivisto.
Questa pagina è stata tradotta automaticamente. Fa fede la versione inglese.

Niente qui è uno schema idealizzato. Ogni estratto proviene così com'è dalla cassa aura-realtime, alla data di pubblicazione. C'è una chiara distinzione qui: NATS JetStream non è ciò che trasporta gli eventi CDC in questa pipeline. Di seguito descriviamo in dettaglio cosa fa effettivamente JetStream qui e cosa trasporta il traffico CDC al suo posto. Più avanti è documentata una trappola Kubernetes verificata, in grado di rendere silenzioso l'intero meccanismo senza generare errori.

L'essenziale
  • Il CDC Postgres di Aurabase legge il WAL tramite wal2json e pg_logical_slot_get_changes - non il protocollo binario pgoutput, non il bridge Debezium/Kafka Connect.
  • Lo slot logico aura_cdc_slot tollera solo un lettore: aura-realtime-cdc-worker viene eletto tramite un Kubernetes Lease (coordination.k8s.io/v1), con una modalità "sempre leader" escludendo K8.
  • Il fan-out delle repliche ws-front passa attraverso NATS core (semplice pub/sub su aura.realtime.>), non attraverso un flusso JetStream persistente. JetStream, in questo stesso servizio, serve esclusivamente la presenza tra istanze KV.
  • Ogni evento viene ricontrollato in RLS dall'abbonato, subito prima della trasmissione, tramite una richiesta SET LOCAL ROLE reale sulla linea interessata, non un'approssimazione memorizzata nella cache.
  • Insidia verificata nella cronologia del repository: senza il plugin wal2json nell'immagine Postgres Kubernetes, la creazione dello slot fallisce e il CDC rimane silenziosamente inattivo.
#
Scelta del motore

wal2json, non pgoutput, non Debezium

La maggior parte delle pipeline CDC Postgres passano attraverso pgoutput, il protocollo di replica della logica binaria, e quindi attraverso un connettore come Debezium che lo traduce in Kafka. Aurabase salta questo passaggio: il servizio aura-realtime decodifica direttamente il WAL nella replica logica con il plugin wal2json, che produce JSON utilizzabile senza una fase di traduzione intermedia.

cdc/postgres.rs (domande reali)sql
-- Creato una volta se assente, ricreato se WAL è scaduto
SELECT pg_create_logical_replication_slot(
  'aura_cdc_slot', 'wal2json');

-- Sondaggio ripetuto (backoff di 100 ms quando il batch ritorna vuoto)
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');

Tre prerequisiti Postgres, verificati nel codice: il ruolo dell'applicazione deve contenere REPLICATION, wal_level = logical deve essere attivo e max_slot_wal_keep_size deve limitare la conservazione WAL. Senza quest'ultimo limite, un consumatore lento fa sì che il volume del disco cresca indefinitamente. Il servizio monitora anche wal_status dello slot: se cambia in lost (WAL eliminato oltre questo limite), lo slot viene ricreato automaticamente. La perdita degli eventi non consumati nel frattempo viene assunta e registrata.

Un batch di 1000 modifiche per sondaggio, con un filtro dinamico della tabella aggiornato ogni 3 secondi. Solo le tabelle in cui un progetto ha abilitato esplicitamente il tempo reale immettono add-tables: ciò evita di decodificare il WAL di tabelle che non interessano gli abbonati.

#
Un unico produttore

cdc-worker: eletto da un Kubernetes Lease, non duplicato

Uno slot di replica logica Postgres tollera solo un'unità attiva alla volta. L'esecuzione di più consumatori dello stesso aura_cdc_slot in parallelo interromperebbe l'ordine delle modifiche, non solo le duplicherebbe. Aurabase ha diviso aura-realtime in due binari separati per risolvere questo vincolo senza sacrificare la scala orizzontale del WebSocket.

Cargo.tomltoml
# cdc-worker: leader eletto, lancia il CDC + pubblica su NATS. Nessun server WS.
[[bin]]
name = "aura-realtime-cdc-worker"

# ws-front: N repliche (HPA), utilizza NATS, filtro RLS, serve WS/SSE.
[[bin]]
name = "aura-realtime-ws-front"

cdc-worker ottiene il diritto di leggere lo slot solo tenendo un oggetto Lease Kubernetes (coordination.k8s.io/v1). Tenta di crearlo, o di rubarlo se è scaduto, quindi lo rinnova dopo un terzo della sua vita. Quando il contratto di locazione viene perso (rinnovo fallito, titolare cambiato), un CancellationToken interrompe il lavoro in corso e il processo ritorna a un ciclo di acquisizione. Al di fuori di Kubernetes, ad esempio durante i test locali o lo sviluppo non cluster, il client K8s non riesce a connettersi e il servizio passa alla modalità “sempre leader”. Utile in fase di sviluppo, errato se lo dimentichi in fase di produzione con diverse repliche.

cdc/leader.rs (firme reali)rust
pub async fn run_with_lease<F, Fut>(cfg: LeaseConfig, work: F) -> anyhow::Result<()>
where F: Fn(CancellationToken) -> Fut, ...
{
  // rinnova tutti lease_duration_secs / 3
  // esclusi K8: funziona (token) solo una volta, token mai annullato
}
#
Un vero e proprio fan-out

Nucleo NATS per CDC, JetStream per presenza

Questo è il punto più spesso frainteso in questo tipo di pipeline: NATS e JetStream sono due cose separate e gli eventi CDC non passano attraverso un flusso JetStream persistente. cdc-worker pubblica ogni evento con Client::publish_with_headers — l'API pub/sub principale NATS, quello consegnato al massimo una volta senza persistenza o riproduzione, separato da JetStream.

Cuore NATS (pub/sub)Fan-out degli eventi CDC alle repliche ws-frontAl massimo una volta, senza persistenza o ripetizione
NATS JetStream (KV)Presenza tra istanze (aura_presence)Stato condiviso anni '60, non un flusso CDC da riprodurre
channels/manager.rs → nats/mod.rs (extraits réels)rust
// cdc-worker pubblica (ChannelManager::publish_and_relay)
nc.publish_with_headers(subject, headers, payload).await

// ws-front ascolta l'intero sottoalbero (NatsRelay::start_relay)
client.subscribe("aura.realtime.>").await

Ogni replica ws-front si iscrive al sottoalbero aura.realtime.>, decodifica l'argomento nel nome del canale originale. Quindi ripubblica l'evento nel suo locale tokio::sync::broadcast, per i WebSocket ad esso collegati. Un'intestazione Aura-Origin porta l'identificatore dell'istanza mittente: ogni replica ignora i messaggi che essa stessa ha pubblicato, evitando così duplicati senza coordinamento centrale.

Questa scelta ha una conseguenza diretta: il core NATS consegna al meglio una volta (at-most-once), senza una coda di lunga durata. Se una replica ws-front viene brevemente disconnessa da NATS al verificarsi di un evento, non verrà recuperata. La cronologia per canale (Channel.history, un buffer in memoria delimitato da HISTORY_SIZE) risiede localmente in ciascuna replica, non a livello di cluster. Questo è un compromesso accettabile proprio perché lo slot di replica logica di Postgres rimane la fonte di verità duratura a monte. È wal2json che garantisce che nessuna modifica al DB venga persa prima del consumo, non NATS.

JetStream esiste in aura-realtime, ma per un uso completamente diverso. Serve l'archivio KV di presenza tra istanze (aura_presence, archivio di memoria, 60 secondi max_age), che sincronizza chi è online su quale canale tra tutte le repliche ws-front. Uno stato condiviso di breve durata, non un flusso di eventi CDC da riprodurre.

#
Sicurezza

RLS ricontrollato dall'abbonato, ad ogni evento

Un evento CDC non viene trasmesso in modo grezzo a tutti gli abbonati di un canale. Il poller instrada ciascuna modifica a uno dei CDC_NUM_SHARDS lavoratori (4 per impostazione predefinita, hash di project_id), che per primo esegue il polling pg_policies. Una tabella senza policy RLS consente a tutti i sottoscrittori (il comportamento Supabase) mentre una tabella con policy attiva un controllo per sottoscrittore.

subscriptions.rs::check_rls_visible_batch (reale, semplificato)sql
BEGIN;
SET LOCAL ROLE "aura_authenticated";
-- Affermazioni JWT dell'abbonato immesse nel GUC
SELECT set_config('request.jwt.claims', $1, true);
SELECT EXISTS(SELECT 1 FROM "schema"."table" WHERE id::text = $2);
ROLLBACK;

La transazione assume il ruolo fisico Postgres dell'abbonato e inserisce le sue attestazioni JWT in request.jwt.claims, quello che auth.uid() e auth.role() leggono dal lato della policy. Poi controlla che la linea rimanga visibile sotto questo ruolo, poi cancella: nessuna scrittura, una vera RLS letta. Solo gli sub_id convalidati atterrano nella allowed_sub_ids dell'evento; un vec vuoto significa nessuno, in fail-closed.

Astuce

Un DELETE sfugge a questo controllo per riga - la riga è scomparsa, è impossibile rileggerla sotto RLS - e ritorna a tutti gli abbonati del canale. In cambio, il payload di un DELETE espone sempre e solo le colonne Identity (chiave primaria), mai il contenuto della riga eliminata.

Un dettaglio che spesso ostacola la migrazione: wal2json cattura il contenuto completo di un UPDATE o di un DELETE solo se la tabella ha REPLICA IDENTITY FULL. Senza di esso, solo INSERT tornano in funzione. Questo è il motivo per cui l'endpoint che abilita il tempo reale su una tabella imposta automaticamente questo ALTER TABLE: una correzione effettiva per il repository, non una casella di controllo separata.

#
Trappola verificata

Perché CDC può rimanere in silenzio in Kubernetes

La storia del fascicolo documenta una trappola reale, non un caso teorico. L'immagine Postgres utilizzata per impostazione predefinita su un cluster Kubernetes (l'immagine della community stock pgvector/pgvector) non include il plug-in wal2json. Senza di esso, pg_create_logical_replication_slot('aura_cdc_slot', 'wal2json') fallisce e nulla sul lato client lo segnala esplicitamente. I WebSocket rimangono aperti, gli abbonamenti vengono accettati, ma non si verificano mai eventi DB.

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

La soluzione effettiva era creare e pubblicare un'immagine Postgres dedicata con questo pacchetto installato, quindi puntare Kubernetes StatefulSet su di essa anziché sull'immagine stock. Inoltre imposta esplicitamente wal_level=logical, max_replication_slots e max_slot_wal_keep_size come argomenti di avvio, assenti dalle immagini generiche.

È proprio per rendere osservabile questo tipo di guasto che cdc-worker espone un flag cdc_active (falso finché lo slot non è confermato integro) sul suo endpoint dedicato /health. Questo è un segnale binario a cui prestare attenzione, piuttosto che dedurre un CDC morto solo dall'assenza di eventi lato client.

#
Limite attuale

Ciò che questa pipeline non copre ancora: cluster dedicati

Questo meccanismo legge una singola variabile POSTGRES_REPLICATION_URL, quindi un singolo server Postgres e un singolo slot logico. Ciò è coerente con lo schema multi-tenant per modello di progetto (project_<uuid>) su un cluster condiviso. Su questo modello tutto funziona: un singolo cdc-worker vede i cambiamenti di tutti i progetti nello stesso cluster e li instrada per diagramma.

Per un progetto su un cluster CNPG dedicato, ovvero la propria istanza Postgres isolata, l'immagine utilizzata oggi (creata dall'immagine stock CloudNativePG) aggiunge solo pg_graphql, non wal2json. Niente, inoltre, esegue un cdc-worker per cluster dedicato. Allo stato attuale del codice, il CDC in tempo reale qui descritto rimane quindi focalizzato sul livello condiviso: un vero e proprio confine architetturale, non una funzionalità della tabella di marcia.

#
Lato cliente

Cosa non cambia: l'API postgres_changes dell'SDK

L'intero meccanismo rimane invisibile dall'SDK. L'API concatenabileaurabase-js non cambia a seconda che l'evento provenga da cdc-worker tramite NATS o, nello sviluppo locale, dal binario unsplit aura-realtime.

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

Ma prima che possa verificarsi un singolo evento, la tabella deve essere registrata presso il CDC: ciò non avviene automaticamente al momento della creazione. Una chiamata service_role sull'endpoint dedicato (o l'interruttore equivalente in Studio) scrive la tabella in realtime_tables dello schema della piattaforma del progetto e imposta REPLICA IDENTITY FULL per te.

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}'

Per il resto della superficie SDK compatibile con Supabase (autenticazione, archiviazione, policy RLS), consultare la guida alla migrazione . Per informazioni su come aura-realtime si divide in più file binari in un unico contenitore dell'area di lavoro, vedere i dettagli dell'architettura Cargo.

PRONTO PER L'IMPLEMENTAZIONE?

Il tuo backend in cinque minuti.

Nessuna carta di credito richiesta · 500 MB gratuiti · 50.000 MAU