PRODSoeverein Europees BaaS-platformOpen Dashboard →

Techniek · 10 min gelezen

PostgreSQL CDC met NATS: real-time architectuur

Affane Daylami · Fondateur · 24 juli 2026

Terug naar blog

Een Aurabase-client die via WebSocket is verbonden, ziet een regel die in de basis is ingevoegd een paar milliseconden na de COMMIT verschijnen – zonder polling, zonder dat er een webhook hoeft te worden geconfigureerd. Het mechanisme bestaat uit drie fasen. Een plug-in voor logische replicatie decodeert de Postgres WAL in JSON, slechts één proces heeft het recht om deze te lezen, en NATS stuurt elke gebeurtenis door naar alle WebSocket-servers in het cluster. Hier ziet u hoe het is aangesloten op de daadwerkelijke code van de aura-realtimeservice.

Deze Engelse tekst is automatisch gegenereerd op basis van het Franse origineel en is nog niet beoordeeld.
Deze pagina is automatisch vertaald. De Engelse versie is gezaghebbend.

Niets hier is een geïdealiseerd patroon. Elk uittreksel komt zoals het is uit de aura-realtime-krat, op de datum van publicatie. Er is hier een duidelijk onderscheid: NATS JetStream transporteert niet CDC-gebeurtenissen in deze pijplijn. Hieronder leggen we uit wat JetStream hier feitelijk doet, en wat daarvoor in de plaats CDC-verkeer vervoert. Een geverifieerde Kubernetes-trap, die in staat is om het hele mechanisme stil te maken zonder fouten te veroorzaken, wordt hieronder verder gedocumenteerd.

De essentie
  • Aurabase's Postgres CDC leest de WAL via wal2json en pg_logical_slot_get_changes – niet het binaire protocol pgoutput, niet de Debezium/Kafka Connect-bridge.
  • Het logische slot aura_cdc_slot tolereert slechts één lezer: aura-realtime-cdc-worker wordt gekozen via een Kubernetes Lease (coordination.k8s.io/v1), met een "altijd leider"-modus exclusief K8's.
  • De fan-out naar de ws-front replica's gaat via NATS core (eenvoudige pub/sub op aura.realtime.>), niet via een aanhoudende JetStream-stream. JetStream bedient in dezelfde service uitsluitend cross-instance-aanwezigheid KV.
  • Elke gebeurtenis wordt vlak voor verzending opnieuw gecontroleerd in RLS door de abonnee, via een echt SET LOCAL ROLE-verzoek op de betreffende lijn - geen in de cache opgeslagen benadering.
  • Valkuil in de geschiedenis van de repository gecontroleerd: zonder de plug-in wal2json in de Postgres Kubernetes-image mislukt het maken van slots en blijft de CDC stil inactief.
#
Motor keuze

wal2json, niet pgoutput, niet Debezium

De meeste CDC Postgres-pijplijnen gaan via pgoutput, het binaire logica-replicatieprotocol, en vervolgens via een connector zoals Debezium die het vertaalt naar Kafka. Aurabase slaat deze stap over: de aura-realtime-service decodeert de WAL direct in logische replicatie met de wal2json-plug-in, die bruikbare JSON produceert zonder tussenliggende vertaalfase.

cdc/postgres.rs (echte vragen)sql
-- Eenmalig gemaakt indien afwezig, opnieuw gemaakt als WAL is verlopen
SELECT pg_create_logical_replication_slot(
  'aura_cdc_slot', 'wal2json');

-- Herhaalde peiling (100 ms uitstel wanneer de batch leeg terugkeert)
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');

Drie Postgres-vereisten, geverifieerd in de code: de applicatierol moet REPLICATIONbevatten, wal_level = logical moet actief zijn en max_slot_wal_keep_size moet WAL-retentie beperken. Zonder deze laatste limiet zorgt een langzame consument ervoor dat het schijfvolume voor onbepaalde tijd groeit. De service bewaakt ook wal_status van het slot: als het verandert in lost (WAL is boven deze limiet verwijderd), wordt het slot automatisch opnieuw aangemaakt. Het verlies van gebeurtenissen die in de tussentijd niet zijn verbruikt, wordt verondersteld en geregistreerd.

Een batch van 1000 wijzigingen per poll, waarbij een dynamisch tabelfilter elke 3 seconden wordt vernieuwd. Alleen tabellen waarvoor een project expliciet real-time heeft ingeschakeld, voeren add-tables in. Dit voorkomt het decoderen van de WAL van tabellen die voor geen enkele abonnee interessant zijn.

#
Eén enkele producent

cdc-worker: gekozen door een Kubernetes Lease, niet gedupliceerd

Een logische replicatiesleuf van Postgres tolereert slechts één actieve schijf tegelijk. Het parallel uitvoeren van meerdere consumenten van dezelfde aura_cdc_slot zou de volgorde van de wijzigingen doorbreken en niet alleen dupliceren. Aurabase splitste aura-realtime in twee afzonderlijke binaire bestanden om deze beperking op te lossen zonder de horizontale schaal van de WebSocket op te offeren.

Cargo.tomltoml
# cdc-worker: gekozen leider, lanceert de CDC + publiceert op NATS. Geen WS-server.
[[bin]]
name = "aura-realtime-cdc-worker"

# ws-front: N replica's (HPA), verbruikt NATS, RLS-filter, bedient WS/SSE.
[[bin]]
name = "aura-realtime-ws-front"

cdc-worker krijgt alleen het recht om het slot te lezen door een Lease Kubernetes-object (coordination.k8s.io/v1) vast te houden. Hij probeert het te maken, of te stelen als het is verlopen, en vernieuwt het vervolgens na een derde van zijn levensduur. Wanneer de huurovereenkomst verloren gaat (verlenging mislukt, houder gewijzigd), beëindigt een CancellationToken het onderhanden werk en keert het proces terug naar een acquisitielus. Buiten Kubernetes – bijvoorbeeld tijdens lokaal testen of niet-clusterontwikkeling – kan de K8s-client geen verbinding maken en schakelt de service over naar de ‘always leading’-modus. Handig bij de ontwikkeling, onjuist als je het vergeet bij productie met meerdere replica's.

cdc/leader.rs (echte handtekeningen)rust
pub async fn run_with_lease<F, Fut>(cfg: LeaseConfig, work: F) -> anyhow::Result<()>
where F: Fn(CancellationToken) -> Fut, ...
{
  // vernieuwt alle lease_duration_secs / 3
  // exclusief K8s: werk(token) slechts één keer, token nooit geannuleerd
}
#
Echte fan-out

NATS-kern voor CDC, JetStream voor aanwezigheid

Dit is het vaakst verkeerd begrepen punt in dit soort pijplijn: NATS en JetStream zijn twee afzonderlijke dingen, en CDC-gebeurtenissen gaan niet door een aanhoudende JetStream-stream. cdc-worker publiceert elke gebeurtenis met Client::publish_with_headers — de NATS core pub/sub API, degene die maximaal één keer levert zonder persistentie of herhaling, los van JetStream.

NATS-hart (pub/sub)Fan-out van CDC-evenementen naar replica's van het ws-frontHoogstens één keer, zonder volharding of herhaling
NATS JetStream (KV)Aanwezigheid tussen meerdere instanties (aura_presence)Gedeelde staat jaren 60, geen CDC-stream om opnieuw af te spelen
channels/manager.rs → nats/mod.rs (extraits réels)rust
// cdc-worker publiceert (ChannelManager::publish_and_relay)
nc.publish_with_headers(subject, headers, payload).await

// ws-front luistert naar de gehele subboom (NatsRelay::start_relay)
client.subscribe("aura.realtime.>").await

Elke replica ws-front abonneert zich op subboom aura.realtime.>en decodeert het onderwerp naar de originele kanaalnaam. Vervolgens publiceert het de gebeurtenis opnieuw naar de lokale tokio::sync::broadcast, voor de WebSockets die ermee zijn verbonden. Een Aura-Origin header draagt ​​de identificatie van de verzendende instantie: elke replica negeert de berichten die hij zelf heeft gepubliceerd, waardoor duplicaten zonder centrale coördinatie worden vermeden.

Deze keuze heeft een direct gevolg: de NATS-kern levert op zijn best één keer (maximaal één keer), zonder een langdurige wachtrij. Als een ws-front-replica kortstondig wordt losgekoppeld van NATS wanneer een gebeurtenis voorbij is, wordt de achterstand niet ingehaald. De geschiedenis per kanaal (Channel.history, een buffer in het geheugen begrensd door HISTORY_SIZE) bevindt zich lokaal op elke replica, niet op clusterniveau. Dit is een aanvaardbaar compromis, juist omdat het logische replicatieslot van Postgres de upstream duurzame bron van waarheid blijft. Het is wal2json dat ervoor zorgt dat er geen DB-wijzigingen verloren gaan vóór consumptie, niet NATS.

JetStream bestaat wel in aura-realtime, maar voor een heel ander gebruik. Het bedient de KV-opslag voor aanwezigheid op meerdere exemplaren (aura_presence, geheugenopslag, 60 seconden max_age), die synchroniseert wie online is op welk kanaal tussen alle ws-front-replica's. Een kortstondige gedeelde staat, geen stroom CDC-gebeurtenissen om opnieuw te spelen.

#
Beveiliging

RLS wordt bij elk evenement opnieuw gecontroleerd door de abonnee

Een CDC-evenement gaat niet onbewerkt naar alle abonnees van een kanaal. De poller stuurt elke wijziging door naar een van de CDC_NUM_SHARDS-werkers (standaard 4, gehasht door project_id), die eerst pg_policiespollt. Een tabel zonder RLS-beleid staat alle abonnees toe (het Supabase-gedrag), terwijl een tabel met beleid een controle per abonnee activeert.

subscriptions.rs::check_rls_visible_batch (echt, vereenvoudigd)sql
BEGIN;
SET LOCAL ROLE "aura_authenticated";
-- JWT-claims van de abonnee geïnjecteerd in GUC
SELECT set_config('request.jwt.claims', $1, true);
SELECT EXISTS(SELECT 1 FROM "schema"."table" WHERE id::text = $2);
ROLLBACK;

De transactie neemt de fysieke Postgres-rol van de abonnee over en injecteert de JWT-claims in request.jwt.claims – degene die auth.uid() en auth.role() aan de beleidszijde lezen. Vervolgens wordt gecontroleerd of de regel zichtbaar blijft onder deze rol en wordt vervolgens geannuleerd: niet geschreven, een echte RLS-lezing. Alleen gevalideerde sub_id komt terecht in allowed_sub_ids van het evenement; een lege vec betekent niemand, in fail-closed.

Astuce

Een DELETE ontsnapt per regel aan deze controle - de regel is verdwenen, herlezen onder RLS is onmogelijk - en gaat terug naar alle abonnees van het kanaal. In ruil daarvoor maakt de payload van een DELETE alleen de identiteitskolommen (primaire sleutel) zichtbaar, nooit de inhoud van de verwijderde rij.

Een detail dat migratie vaak in de weg staat: wal2json legt alleen de volledige inhoud van een UPDATE of een DELETE vast als de tabel REPLICA IDENTITY FULLbevat. Zonder dit komt alleen de INSERT terug. Dit is de reden waarom het eindpunt dat real-time op een tabel mogelijk maakt, deze ALTER TABLE automatisch instelt - een daadwerkelijke oplossing voor de repository, en niet een apart selectievakje.

#
Val geverifieerd

Waarom CDC kan zwijgen in Kubernetes

De geschiedenis van de indiening documenteert een echte valstrik, geen theoretisch geval. De Postgres-image die standaard wordt gebruikt op een Kubernetes-cluster (de pgvector/pgvector stock community-image) bevat niet de plug-in wal2json. Zonder dit mislukt pg_create_logical_replication_slot('aura_cdc_slot', 'wal2json'), en niets aan de clientzijde signaleert dit expliciet. WebSockets blijven open, abonnementen worden geaccepteerd, maar er vinden nooit DB-gebeurtenissen plaats.

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

De feitelijke oplossing was het bouwen en publiceren van een speciale Postgres-image waarop dit pakket was geïnstalleerd, en vervolgens de Kubernetes StatefulSet ernaar te richten in plaats van de stock-image. Het stelde ook expliciet wal_level=logical, max_replication_slots en max_slot_wal_keep_size in als opstartargumenten - afwezig in algemene afbeeldingen.

Het is precies om dit soort fouten waarneembaar te maken dat cdc-worker een cdc_active-vlag (onwaar zolang het slot niet als gezond wordt bevestigd) op het speciale /health-eindpunt weergeeft. Dit is een binair signaal waar we op moeten letten – in plaats van dode CDC af te leiden uit de afwezigheid van gebeurtenissen aan de clientzijde.

#
Huidige limiet

Wat deze pijplijn nog niet omvat: dedicated clusters

Dit mechanisme leest een enkele POSTGRES_REPLICATION_URLvariabele, dus een enkele Postgres-server en een enkel logisch slot. Dit komt overeen met het multi-tenant schema per projectmodel (project_<uuid>) op een gedeeld cluster. Op dit model werkt alles: een enkele cdc-worker ziet de wijzigingen van alle projecten in hetzelfde cluster, en route per diagram.

Voor een project op een speciaal CNPG-cluster (zijn eigen geïsoleerde Postgres-instantie) voegt de afbeelding die vandaag wordt gebruikt (opgebouwd op basis van de stockimage van CloudNativePG) alleen pg_graphqltoe, niet wal2json. Ook niets voert een cdc-worker per speciaal cluster uit. In de huidige staat van de code blijft de hier beschreven real-time CDC daarom gefocust op de gedeelde laag: een echte architectonische grens, geen routekaartfunctionaliteit.

#
Klantzijde

Wat niet verandert: de postgres_changes API van de SDK

Dit hele mechanisme blijft onzichtbaar vanuit de SDK. De aaneenschakelbare APIaurabase-js verandert niet, afhankelijk van of de gebeurtenis afkomstig is van cdc-worker via NATS of, bij lokale ontwikkelaars, van het niet-gesplitste binaire bestand aura-realtime.

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

Maar voordat er ook maar één gebeurtenis kan plaatsvinden, moet de tabel bij de CDC worden geregistreerd; dit gebeurt niet automatisch bij het maken ervan. Een service_role-aanroep op het toegewezen eindpunt (of de equivalente schakelaar in Studio) schrijft de tabel in realtime_tables van het platformschema van het project en stelt REPLICA IDENTITY FULL voor u in.

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

Voor de rest van het Supabase-compatibele SDK-oppervlak (authenticatie, opslag, beleid RLS) raadpleegt u de migratiehandleiding. Voor hoe aura-realtime zich opsplitst in meerdere binaire bestanden in één werkruimtekrat, zie het Cargo-architectuurdetail.

KLAAR VOOR IMPLEMENTATIE?

Uw backend in vijf minuten.

Geen creditcard vereist · 500 MB gratis · 50.000 MAU