Nic tutaj nie jest wyidealizowanym wzorcem. Każdy ekstrakt pochodzi ze skrzynki aura-realtimew dniu publikacji. Istnieje tu wyraźne rozróżnienie: NATS JetStream nie jest tym, co transportuje zdarzenia CDC w tym potoku. Poniżej szczegółowo opisujemy, co faktycznie robi tutaj JetStream i co zamiast tego przenosi ruch CDC. Sprawdzoną pułapkę Kubernetesa, zdolną do wyciszenia całego mechanizmu bez powodowania jakichkolwiek błędów, opisano poniżej.
- CDC Postgres firmy Aurabase odczytuje WAL poprzez
wal2jsonipg_logical_slot_get_changes— a nie protokół binarnypgoutput, a nie most Debezium/Kafka Connect. - Gniazdo logiczne
aura_cdc_slottoleruje tylko jeden czytnik:aura-realtime-cdc-workerjest wybierany za pośrednictwem dzierżawy Kubernetes (coordination.k8s.io/v1), z trybem „zawsze lidera” z wyłączeniem K8. - Rozprzestrzenianie się do replik
ws-frontprzechodzi przez rdzeń NATS (prosty pub/sub naaura.realtime.>), a nie przez trwały strumień JetStream. JetStream w tej samej usłudze obsługuje wyłącznie obecność KV między instancjami. - Każde zdarzenie jest ponownie sprawdzane w RLS przez abonenta, tuż przed transmisją, poprzez rzeczywiste żądanie
SET LOCAL ROLEna danej linii — a nie przybliżenie w pamięci podręcznej. - Sprawdzona pułapka w historii repozytorium: bez wtyczki
wal2jsonw obrazie Postgres Kubernetes tworzenie slotu kończy się niepowodzeniem, a CDC pozostaje cicho nieaktywne.
wal2json, nie pgoutput, nie Debezium
Większość potoków CDC Postgres przechodzi przez pgoutput, protokół replikacji logiki binarnej, a następnie przez złącze takie jak Debezium, które tłumaczy to na Kafkę. Aurabase pomija ten krok: usługa aura-realtime bezpośrednio dekoduje WAL do replikacji logicznej za pomocą wtyczki wal2json, która tworzy użyteczny kod JSON bez pośredniego etapu tłumaczenia.
Trzy wymagania wstępne Postgres zweryfikowane w kodzie: rola aplikacji musi zawierać REPLICATION, wal_level = logical musi być aktywna, a max_slot_wal_keep_size musi ograniczać przechowywanie WAL. Bez tego ostatniego ograniczenia powolny konsument powoduje nieograniczony wzrost woluminu dysku. Usługa monitoruje również wal_status slotu: jeśli zmieni się na lost (WAL został usunięty poza ten limit), slot zostanie automatycznie utworzony ponownie. Zakładana jest i rejestrowana utrata zdarzeń, które w międzyczasie nie zostały wykorzystane.
Partia 1000 zmian na sondę, z dynamicznym filtrem tabeli odświeżanym co 3 sekundy. Tylko tabele, dla których projekt wyraźnie włączył czas rzeczywisty, wpisz add-tables — pozwala to uniknąć dekodowania WAL tabel, które nie są interesujące dla żadnego subskrybenta.
cdc-worker: wybrany w ramach dzierżawy Kubernetes, nie duplikowany
Gniazdo replikacji logicznej Postgres toleruje tylko jeden aktywny dysk na raz. Równoległe uruchomienie wielu konsumentów tego samego aura_cdc_slot spowodowałoby naruszenie kolejności zmian, a nie tylko ich zduplikowanie. Aurabase podzielił aura-realtime na dwa oddzielne pliki binarne, aby rozwiązać to ograniczenie bez poświęcania poziomej skali protokołu WebSocket.
cdc-worker uzyskuje prawo do odczytu slotu jedynie poprzez trzymanie obiektu Lease Kubernetes (coordination.k8s.io/v1). Próbuje go stworzyć lub ukraść, jeśli wygasł, a następnie odnawia go po jednej trzeciej jego żywotności. W przypadku utraty dzierżawy (odnowienie nie powiodło się, zmienił się posiadacz), CancellationToken przerywa pracę w toku i proces powraca do pętli przejęcia. Poza Kubernetesem — na przykład podczas testów lokalnych lub programowania poza klastrem — klient K8s nie łączy się, a usługa przełącza się w tryb „zawsze wiodący”. Przydatne w deweloperach, niepoprawne, jeśli zapomnisz o tym w prod z kilkoma replikami.
Rdzeń NATS dla CDC, JetStream dla obecności
Jest to najczęściej źle rozumiany punkt w tego rodzaju potoku: NATS i JetStream to dwie różne rzeczy, a zdarzenia CDC nie przechodzą przez trwały strumień JetStream. cdc-worker publikuje każde zdarzenie za pomocą Client::publish_with_headers — podstawowego API pub/sub NATS, tego, które zostało dostarczone najwyżej raz bez utrwalania i odtwarzania, niezależnie od JetStream.
| Serce NATS (pub/sub) | Rozprzestrzenianie zdarzeń CDC do replik ws-front | Co najwyżej raz, bez utrzymywania się i powtarzania |
|---|---|---|
| NATS JetStream (KV) | Obecność między instancjami (aura_presence) | Udostępniony stan 60., a nie strumień CDC do odtworzenia |
Każda replika ws-front subskrybuje poddrzewo aura.realtime.>, dekoduje temat do oryginalnej nazwy kanału. Następnie ponownie publikuje zdarzenie w swoim lokalnym tokio::sync::broadcastdla podłączonych do niego gniazd WebSocket. Nagłówek Aura-Origin zawiera identyfikator instancji wysyłającej: każda replika ignoruje komunikaty, które sama opublikowała, co pozwala uniknąć duplikatów bez centralnej koordynacji.
Wybór ten ma bezpośrednie konsekwencje: rdzeń NATS dostarcza co najwyżej raz (co najwyżej raz), bez długotrwałej kolejki. Jeśli replika ws-front zostanie na krótko odłączona od NATS po zakończeniu zdarzenia, nie nadrobi zaległości. Historia na kanał (Channel.history, bufor w pamięci ograniczony przez HISTORY_SIZE) żyje lokalnie w każdej replice, a nie na poziomie klastra. Jest to akceptowalny kompromis właśnie dlatego, że gniazdo replikacji logicznej Postgres pozostaje trwałym źródłem prawdy. To wal2json zapewnia, że żadne zmiany DB nie zostaną utracone przed zużyciem, a nie NATS.
JetStream istnieje w aura-realtime — ale do zupełnie innego zastosowania. Obsługuje magazyn KV obecności między instancjami (aura_presence, magazyn pamięci, 60-sekundowy max_age), który synchronizuje, kto jest online na jakim kanale pomiędzy wszystkimi replikami ws-front. Krótkotrwały stan współdzielony, a nie strumień zdarzeń CDC do odtworzenia.
RLS jest ponownie sprawdzany przez abonenta przy każdym zdarzeniu
Zdarzenie CDC nie jest dostępne dla wszystkich subskrybentów kanału. Moduł odpytujący kieruje każdą zmianę do jednego z procesów roboczych CDC_NUM_SHARDS (domyślnie 4, szyfrowane przez project_id), który najpierw odpytuje pg_policies. Tabela bez zasad RLS umożliwia wszystkim subskrybentom — zachowanie Supabase — podczas gdy tabela z zasadami uruchamia kontrolę dla każdego subskrybenta.
Transakcja przejmuje fizyczną rolę subskrybenta w Postgres i wprowadza jego roszczenia JWT do request.jwt.claims — tego, który auth.uid() i auth.role() czytają po stronie polityki. Następnie sprawdza, czy linia pozostaje widoczna w tej roli, po czym anuluje: brak zapisu, prawdziwy odczyt RLS. Tylko zatwierdzone sub_id lądują w allowed_sub_ids wydarzenia; pusty vec oznacza nikogo, w zamkniętym awaryjnie.
DELETE omija tę kontrolę w każdym wierszu - wiersz zniknął, nie można go ponownie odczytać w ramach RLS - i wraca do wszystkich abonentów kanału. W zamian ładunek DELETE zawsze ujawnia tylko kolumny tożsamości (klucz podstawowy), a nigdy zawartość usuniętego wiersza.
Szczegół, który często utrudnia migrację: wal2json przechwytuje tylko całą zawartość UPDATE lub DELETE, jeśli tabela zawiera REPLICA IDENTITY FULL. Bez tego pojawią się tylko INSERT. Właśnie dlatego punkt końcowy, który umożliwia pracę w czasie rzeczywistym na tabeli, automatycznie ustawia to ALTER TABLE — rzeczywistą poprawkę do repozytorium, a nie osobne pole wyboru.
Dlaczego CDC może milczeć w Kubernetesie
Historia zgłoszenia dokumentuje prawdziwą pułapkę, a nie teoretyczny przypadek. Obraz Postgres używany domyślnie w klastrze Kubernetes (obraz społecznościowy pgvector/pgvector) nie zawiera wtyczki wal2json. Bez tego pg_create_logical_replication_slot('aura_cdc_slot', 'wal2json') kończy się niepowodzeniem i nic po stronie klienta wyraźnie tego nie sygnalizuje. WebSockets pozostają otwarte, subskrypcje są akceptowane, ale nie zachodzą żadne zdarzenia DB.
Faktyczną poprawką było zbudowanie i opublikowanie dedykowanego obrazu Postgres z zainstalowanym tym pakietem, a następnie skierowanie na niego Kubernetes StatefulSet zamiast na obraz podstawowy. Wyraźnie ustawia również wal_level=logical, max_replication_slots i max_slot_wal_keep_size jako argumenty startowe - nieobecne w obrazach ogólnych.
Właśnie po to, aby tego typu awarie były zauważalne, cdc-worker udostępnia flagę cdc_active (fałszywą, o ile gniazdo nie jest potwierdzone, że jest zdrowe) w dedykowanym punkcie końcowym /health. Jest to sygnał binarny, na który należy zwrócić uwagę — zamiast wnioskować o martwym CDC na podstawie samego braku zdarzeń po stronie klienta.
Czego ten potok jeszcze nie obejmuje: dedykowane klastry
Mechanizm ten odczytuje pojedynczą zmienną POSTGRES_REPLICATION_URL, a zatem pojedynczy serwer Postgres i pojedynczy slot logiczny. Jest to spójne ze schematem wielu dzierżawców na model projektu (project_<uuid>) w klastrze udostępnionym. W tym modelu wszystko działa: pojedynczy cdc-worker widzi zmiany we wszystkich projektach w tym samym klastrze i trasę według diagramu.
W przypadku projektu w dedykowanym klastrze CNPG — własnej, izolowanej instancji Postgres — używany dzisiaj obraz (zbudowany z podstawowego obrazu CloudNativePG) dodaje tylko pg_graphql, a nie wal2json. Nic też nie uruchamia cdc-worker na dedykowany klaster. Dlatego też w obecnym stanie kodu opisane tutaj CDC działające w czasie rzeczywistym pozostaje skupione na warstwie współdzielonej — prawdziwej granicy architektonicznej, a nie funkcjonalności planu działania.
Co się nie zmienia: interfejs API postgres_changes pakietu SDK
Cały ten mechanizm pozostaje niewidoczny z SDK. Łańcuchowy interfejs APIaurabase-js nie zmienia się w zależności od tego, czy zdarzenie pochodzi z cdc-worker przez NATS, czy w przypadku lokalnego dewelopera z niedzielonego pliku binarnego aura-realtime.
Zanim jednak może nastąpić pojedyncze zdarzenie, tabela musi zostać zarejestrowana w CDC — nie następuje to automatycznie w momencie utworzenia. Wywołanie service_role na dedykowanym punkcie końcowym (lub równoważny przełącznik w Studio) zapisuje tabelę w realtime_tables schematu platformy projektu i ustawia REPLICA IDENTITY FULL.
Aby poznać resztę pakietu SDK zgodnego z Supabase — uwierzytelnianie, przechowywanie, zasady RLS — zobacz przewodnik migracji . Aby dowiedzieć się, jak aura-realtime dzieli się na wiele plików binarnych w jednej skrzyni obszaru roboczego, zobacz szczegóły architektury Cargo.