PRODSouveräne europäische BaaS-PlattformÖffnen Sie das Dashboard →

Ingenieurwesen · 10 Min. Lesezeit

PostgreSQL CDC mit NATS: Echtzeitarchitektur

Affane Daylami · Fondateur · 24. Juli 2026

Zurück zum Blog

Ein über WebSocket verbundener Aurabase-Client sieht, dass einige Millisekunden nach dem COMMIT eine in die Basis eingefügte Zeile erscheint – ohne Abfrage, ohne zu konfigurierenden Webhook. Der Mechanismus ist dreistufig. Ein logisches Replikations-Plugin dekodiert die Postgres-WAL in JSON, nur ein Prozess hat das Recht, sie zu lesen, und NATS leitet jedes Ereignis an alle WebSocket-Server im Cluster weiter. So ist es im eigentlichen Code des Aura-Realtimeservice verkabelt.

Dieser englische Text wurde automatisch aus dem französischen Original generiert und wurde noch nicht überprüft.
Diese Seite wurde automatisch übersetzt. Maßgeblich ist die englische Version.

Nichts hier ist ein idealisiertes Muster. Jeder Auszug stammt zum Zeitpunkt der Veröffentlichung unverändert aus der Kiste aura-realtime. Hier gibt es eine klare Unterscheidung: NATS JetStream ist nicht das, was CDC-Ereignisse in dieser Pipeline transportiert. Im Folgenden erläutern wir detailliert, was JetStream hier eigentlich macht und was stattdessen den CDC-Verkehr überträgt. Eine verifizierte Kubernetes-Falle, die in der Lage ist, den gesamten Mechanismus stumm zu schalten, ohne Fehler auszulösen, ist weiter unten dokumentiert.

Das Wesentliche
  • Der Postgres CDC von Aurabase liest die WAL über wal2json und pg_logical_slot_get_changes – nicht das pgoutput-Binärprotokoll, nicht die Debezium/Kafka Connect-Brücke.
  • Der logische Steckplatz aura_cdc_slot toleriert nur einen Leser: aura-realtime-cdc-worker wird über einen Kubernetes-Lease (coordination.k8s.io/v1) ausgewählt, mit einem „Always Leader“-Modus ohne K8s.
  • Der Fanout zu den ws-front-Replikaten erfolgt über den NATS -Kern (einfaches Pub/Sub auf aura.realtime.>), nicht über einen dauerhaften JetStream-Stream. JetStream bedient im selben Dienst ausschließlich instanzübergreifende Präsenz-KV.
  • Jedes Ereignis wird in RLS vom Abonnenten kurz vor der Übertragung erneut überprüft, und zwar über eine echte SET LOCAL ROLE-Anfrage auf der betreffenden Leitung – keine zwischengespeicherte Annäherung.
  • Überprüfter Fallstrick im Repository-Verlauf: Ohne das wal2json-Plugin im Postgres-Kubernetes-Image schlägt die Slot-Erstellung fehl und der CDC bleibt stillschweigend inaktiv.
#
Motorauswahl

wal2json, nicht pgoutput, nicht Debezium

Die meisten CDC-Postgres-Pipelines durchlaufen pgoutput, das binäre Logik-Replikationsprotokoll, und dann einen Connector wie Debezium, der es in Kafka übersetzt. Aurabase überspringt diesen Schritt: Der aura-realtime-Dienst dekodiert die WAL direkt in die logische Replikation mit dem Wal2json-Plugin, das ohne Zwischenübersetzungsstufe verwendbares JSON erzeugt.

cdc/postgres.rs (echte Fragen)sql
-- Wird bei Abwesenheit einmal erstellt und bei Ablauf der WAL neu erstellt
SELECT pg_create_logical_replication_slot(
  'aura_cdc_slot', 'wal2json');

-- Wiederholte Abfrage (Backoff 100 ms, wenn der Stapel leer zurückkommt)
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');

Drei Postgres-Voraussetzungen, die im Code überprüft werden: Die Anwendungsrolle muss REPLICATIONtragen, wal_level = logical muss aktiv sein und max_slot_wal_keep_size muss die WAL-Aufbewahrung begrenzen. Ohne diese letzte Grenze führt ein langsamer Verbraucher dazu, dass das Festplattenvolumen auf unbestimmte Zeit anwächst. Der Dienst überwacht auch wal_status des Steckplatzes: Wenn er sich in lost ändert (WAL über diese Grenze hinaus gelöscht), wird der Steckplatz automatisch neu erstellt. Der Verlust zwischenzeitlich nicht konsumierter Ereignisse wird angenommen und protokolliert.

Ein Stapel von 1000 Änderungen pro Umfrage, wobei alle 3 Sekunden ein dynamischer Tabellenfilter aktualisiert wird. Nur Tabellen, bei denen ein Projekt explizit Echtzeit aktiviert hat, geben add-tables ein – dadurch wird die Dekodierung der WAL von Tabellen vermieden, die für keinen Abonnenten von Interesse sind.

#
Ein einziger Produzent

cdc-worker: durch einen Kubernetes-Lease gewählt, nicht dupliziert

Ein logischer Postgres-Replikationssteckplatz toleriert jeweils nur ein aktives Laufwerk. Das parallele Ausführen mehrerer Verbraucher desselben aura_cdc_slot würde die Reihenfolge der Änderungen zerstören und sie nicht nur duplizieren. Aurabase hat aura-realtime in zwei separate Binärdateien aufgeteilt, um diese Einschränkung zu lösen, ohne die horizontale Skalierung des WebSocket zu beeinträchtigen.

Cargo.tomltoml
# cdc-worker: gewählter Anführer, startet die CDC + veröffentlicht auf NATS. Kein WS-Server.
[[bin]]
name = "aura-realtime-cdc-worker"

# ws-front: N Replikate (HPA), nutzt NATS, RLS-Filter, bedient WS/SSE.
[[bin]]
name = "aura-realtime-ws-front"

cdc-worker erhält nur das Recht, den Slot zu lesen, indem es ein Lease Kubernetes-Objekt (coordination.k8s.io/v1) hält. Er versucht, es zu erschaffen oder es zu stehlen, wenn es abgelaufen ist, und erneuert es dann nach einem Drittel seiner Lebensdauer. Wenn der Mietvertrag verloren geht (Verlängerung fehlgeschlagen, Inhaber geändert), unterbricht ein CancellationToken die laufende Arbeit und der Prozess kehrt zu einer Erfassungsschleife zurück. Außerhalb von Kubernetes – beispielsweise während lokaler Tests oder Nicht-Cluster-Entwicklung – kann der K8s-Client keine Verbindung herstellen und der Dienst wechselt in den „Always Leading“-Modus. Nützlich in der Entwicklung, falsch, wenn Sie es in Produkten mit mehreren Replikaten vergessen.

cdc/leader.rs (echte Unterschriften)rust
pub async fn run_with_lease<F, Fut>(cfg: LeaseConfig, work: F) -> anyhow::Result<()>
where F: Fn(CancellationToken) -> Fut, ...
{
  // erneuert alle lease_duration_secs / 3
  // ausgenommen K8s: Arbeit (Token) nur einmal, Token nie storniert
}
#
Echter Fan-Out

NATS-Kern für CDC, JetStream für Präsenz

Dies ist der am häufigsten missverstandene Punkt bei dieser Art von Pipeline: NATS und JetStream sind zwei getrennte Dinge, und CDC-Ereignisse durchlaufen keinen dauerhaften JetStream-Stream. cdc-worker veröffentlicht jedes Ereignis mit Client::publish_with_headers – der NATS-Kern-Pub/Sub-API, der höchstens einmal ohne Persistenz oder Wiedergabe, getrennt von JetStream.

NATS-Herz (Pub/Sub)Weitergabe von CDC-Ereignissen an WS-Front-ReplikateHöchstens einmal, ohne Persistenz oder Wiederholung
NATS JetStream (KV)Instanzübergreifende Präsenz (aura_presence)Gemeinsamer Status 60s, kein CDC-Stream zum Wiedergeben
channels/manager.rs → nats/mod.rs (extraits réels)rust
// cdc-worker veröffentlicht (ChannelManager::publish_and_relay)
nc.publish_with_headers(subject, headers, payload).await

// ws-front lauscht auf den gesamten Teilbaum (NatsRelay::start_relay)
client.subscribe("aura.realtime.>").await

Jedes Replikat ws-front abonniert den Teilbaum aura.realtime.>und dekodiert das Thema in den ursprünglichen Kanalnamen. Anschließend veröffentlicht es das Ereignis erneut in seinem lokalen tokio::sync::broadcastfür die damit verbundenen WebSockets. Ein Aura-Origin-Header trägt die Kennung der sendenden Instanz: Jedes Replikat ignoriert die Nachrichten, die es selbst veröffentlicht hat, wodurch Duplikate ohne zentrale Koordination vermieden werden.

Diese Wahl hat eine direkte Konsequenz: Der NATS-Kern liefert höchstens einmal (at-most-once), ohne eine lange Warteschlange. Wenn ein ws-front-Replikat bei Ablauf eines Ereignisses kurzzeitig von NATS getrennt wird, kann es nicht aufholen. Der Verlauf pro Kanal (Channel.history, ein durch HISTORY_SIZEbegrenzter In-Memory-Puffer) lebt lokal auf jedem Replikat, nicht auf Clusterebene. Dies ist ein akzeptabler Kompromiss, gerade weil der logische Replikationsslot von Postgres die vorgelagerte dauerhafte Quelle der Wahrheit bleibt. Es ist wal2json, das sicherstellt, dass keine DB-Änderungen vor dem Verbrauch verloren gehen, nicht NATS.

JetStream existiert zwar in aura-realtime – aber für einen ganz anderen Zweck. Es bedient den instanzübergreifenden Präsenz-KV-Speicher (aura_presence, Speicherspeicher, 60 Sekunden max_age), der zwischen allen ws-front-Replikaten synchronisiert, wer auf welchem ​​Kanal online ist. Ein kurzlebiger gemeinsamer Zustand, kein Strom von CDC-Ereignissen zum Wiederholen.

#
Sicherheit

RLS wird bei jeder Veranstaltung vom Abonnenten erneut überprüft

Ein CDC-Ereignis wird nicht an alle Abonnenten eines Kanals weitergeleitet. Der Poller leitet jede Änderung an einen der CDC_NUM_SHARDS-Worker weiter (standardmäßig 4, gehasht durch project_id), der zuerst pg_policiesabfragt. Eine Tabelle ohne RLS-Richtlinie lässt alle Abonnenten zu – das Supabase-Verhalten –, während eine Tabelle mit Richtlinien eine Prüfung pro Abonnenten auslöst.

subscriptions.rs::check_rls_visible_batch (real, vereinfacht)sql
BEGIN;
SET LOCAL ROLE "aura_authenticated";
-- JWT-Ansprüche des Abonnenten in GUC eingespeist
SELECT set_config('request.jwt.claims', $1, true);
SELECT EXISTS(SELECT 1 FROM "schema"."table" WHERE id::text = $2);
ROLLBACK;

Die Transaktion übernimmt die physische Postgres-Rolle des Abonnenten und fügt ihre JWT-Ansprüche in request.jwt.claims ein – den, den auth.uid() und auth.role() auf der Richtlinienseite lesen. Anschließend prüft es, ob die Zeile unter dieser Rolle sichtbar bleibt, und bricht dann ab: kein Schreiben, ein echter RLS-Lesevorgang. Nur validierte sub_id landen in allowed_sub_ids des Ereignisses; Ein leeres VEC bedeutet niemand, in fail-closed.

Scharfsinn

Ein DELETE entgeht dieser Prüfung pro Zeile – die Zeile ist verschwunden, es ist unmöglich, sie unter RLS erneut zu lesen – und geht an alle Abonnenten des Kanals zurück. Im Gegenzug legt die Nutzlast eines DELETE immer nur die Identitätsspalten (Primärschlüssel) offen, niemals den Inhalt der gelöschten Zeile.

Ein Detail, das der Migration oft im Weg steht: wal2json erfasst nur dann den vollständigen Inhalt eines UPDATE oder eines DELETE, wenn die Tabelle REPLICA IDENTITY FULLhat. Ohne sie wird nur INSERT wieder angezeigt. Aus diesem Grund setzt der Endpunkt, der Echtzeit für eine Tabelle ermöglicht, diesen ALTER TABLE automatisch – eine tatsächliche Korrektur für das Repository, kein separates Kontrollkästchen.

#
Falle verifiziert

Warum CDC in Kubernetes schweigen kann

Die Aktenhistorie dokumentiert eine echte Falle, keinen theoretischen Fall. Das standardmäßig in einem Kubernetes-Cluster verwendete Postgres-Image (das pgvector/pgvector-Stock-Community-Image) enthält nicht das Plugin wal2json. Ohne sie schlägt pg_create_logical_replication_slot('aura_cdc_slot', 'wal2json') fehl und nichts auf der Clientseite signalisiert dies explizit. WebSockets bleiben geöffnet, Abonnements werden akzeptiert, es finden jedoch nie DB-Ereignisse statt.

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

Die eigentliche Lösung bestand darin, ein dediziertes Postgres-Image mit diesem installierten Paket zu erstellen und zu veröffentlichen und dann den Kubernetes StatefulSet darauf statt auf das Archiv-Image zu verweisen. Außerdem werden explizit wal_level=logical, max_replication_slots und max_slot_wal_keep_size als Startargumente festgelegt – in generischen Bildern nicht vorhanden.

Genau um diese Art von Fehlern beobachtbar zu machen, stellt cdc-worker ein cdc_active-Flag (falsch, solange der Steckplatz nicht als fehlerfrei bestätigt wird) auf seinem dedizierten /health-Endpunkt bereit. Dies ist ein binäres Signal, auf das man achten sollte – und nicht allein aus dem Fehlen clientseitiger Ereignisse auf einen toten CDC schließen.

#
Aktuelle Grenze

Was diese Pipeline noch nicht abdeckt: dedizierte Cluster

Dieser Mechanismus liest eine einzelne POSTGRES_REPLICATION_URL-Variable, also einen einzelnen Postgres-Server und einen einzelnen logischen Steckplatz. Dies steht im Einklang mit dem mehrinstanzenfähigen Schema pro Projektmodell (project_<uuid>) auf einem gemeinsam genutzten Cluster. Bei diesem Modell funktioniert alles: Ein einziger cdc-worker sieht die Änderungen aller Projekte im selben Cluster und leitet sie per Diagramm weiter.

Für ein Projekt auf einem dedizierten CNPG-Cluster – einer eigenen isolierten Postgres-Instanz – fügt das heute verwendete Image (erstellt aus dem CloudNativePG-Archiv-Image) nur pg_graphqlhinzu, nicht wal2json. Auch nichts führt einen cdc-worker pro dediziertem Cluster aus. Im aktuellen Zustand des Codes konzentriert sich der hier beschriebene Echtzeit-CDC daher weiterhin auf die gemeinsame Ebene – eine echte Architekturgrenze, keine Roadmap-Funktionalität.

#
Kundenseite

Was sich nicht ändert: die postgres_changes API des SDK

Dieser gesamte Mechanismus bleibt für das SDK unsichtbar. Die verkettbare APIaurabase-js ändert sich nicht, je nachdem, ob das Ereignis von cdc-worker über NATS oder, in der lokalen Entwicklung, von der ungeteilten Binärdatei aura-realtime kommt.

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

Doch bevor ein einzelnes Ereignis eintreten kann, muss die Tabelle beim CDC registriert werden – dies erfolgt nicht automatisch bei der Erstellung. Ein service_role-Aufruf auf dem dedizierten Endpunkt (oder der entsprechende Schalter in Studio) schreibt die Tabelle in realtime_tables des Plattformschemas des Projekts und legt REPLICA IDENTITY FULL für Sie fest.

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

Weitere Informationen zur Supabase-kompatiblen SDK-Oberfläche – Authentifizierung, Speicher, RLS-Richtlinien – finden Sie im Migrationsleitfaden . Informationen zur Aufteilung von aura-realtime in mehrere Binärdateien in einer einzelnen Arbeitsbereichskiste finden Sie im Detail zur Cargo-Architektur .

BEREIT ZUM EINSATZ?

Ihr Backend in fünf Minuten.

Keine Kreditkarte erforderlich · 500 MB kostenlos · 50.000 MAU