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.
- Der Postgres CDC von Aurabase liest die WAL über
wal2jsonundpg_logical_slot_get_changes– nicht daspgoutput-Binärprotokoll, nicht die Debezium/Kafka Connect-Brücke. - Der logische Steckplatz
aura_cdc_slottoleriert nur einen Leser:aura-realtime-cdc-workerwird ü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 aufaura.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.
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.
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.
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.
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.
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-Replikate | Höchstens einmal, ohne Persistenz oder Wiederholung |
|---|---|---|
| NATS JetStream (KV) | Instanzübergreifende Präsenz (aura_presence) | Gemeinsamer Status 60s, kein CDC-Stream zum Wiedergeben |
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.
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.
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.
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.
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.
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.
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.
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.
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.
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 .