PRODSuwerenna europejska platforma BaaSOtwórz Panel →

Inżynieria · 10 min odczytu

PostgreSQL CDC z NATS: architektura czasu rzeczywistego

Affane Daylami · Fondateur · 24 lipca 2026

Powrót do bloga

Klient Aurabase podłączony przez WebSocket widzi, że linia wstawiona do bazy pojawia się kilka milisekund po COMMIT – bez odpytywania i konfigurowania webhooka. Mechanizm składa się z trzech etapów. Wtyczka replikacji logicznej dekoduje Postgres WAL do formatu JSON, tylko jeden proces ma prawo go odczytać, a NATS przekazuje każde zdarzenie do wszystkich serwerów WebSocket w klastrze. Oto jak jest to podłączone w rzeczywistym kodzie usługi aura-realtimeservice.

Ten tekst w języku angielskim został wygenerowany automatycznie na podstawie francuskiego oryginału i nie był jeszcze recenzowany.
Ta strona została przetłumaczona automatycznie. Wersja angielska jest miarodajna.

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.

Najważniejsze
  • CDC Postgres firmy Aurabase odczytuje WAL poprzez wal2json i pg_logical_slot_get_changes — a nie protokół binarny pgoutput, a nie most Debezium/Kafka Connect.
  • Gniazdo logiczne aura_cdc_slot toleruje tylko jeden czytnik: aura-realtime-cdc-worker jest wybierany za pośrednictwem dzierżawy Kubernetes (coordination.k8s.io/v1), z trybem „zawsze lidera” z wyłączeniem K8.
  • Rozprzestrzenianie się do replik ws-front przechodzi przez rdzeń NATS (prosty pub/sub na aura.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 ROLE na danej linii — a nie przybliżenie w pamięci podręcznej.
  • Sprawdzona pułapka w historii repozytorium: bez wtyczki wal2json w obrazie Postgres Kubernetes tworzenie slotu kończy się niepowodzeniem, a CDC pozostaje cicho nieaktywne.
#
Wybór silnika

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.

cdc/postgres.rs (prawdziwe zapytania)sql
-- Utworzony raz w przypadku braku, odtworzony w przypadku wygaśnięcia WAL
SELECT pg_create_logical_replication_slot(
  'aura_cdc_slot', 'wal2json');

-- Powtarzane odpytywanie (cofnięcie 100 ms, gdy partia wraca pusta)
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');

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.

#
Jeden producent

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.

Cargo.tomltoml
# cdc-worker: wybrany lider, uruchamia CDC + publikuje na NATS. Brak serwera WS.
[[bin]]
name = "aura-realtime-cdc-worker"

# ws-front: N replik (HPA), zużywa NATS, filtr RLS, obsługuje WS/SSE.
[[bin]]
name = "aura-realtime-ws-front"

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.

cdc/leader.rs (prawdziwe podpisy)rust
pub async fn run_with_lease<F, Fut>(cfg: LeaseConfig, work: F) -> anyhow::Result<()>
where F: Fn(CancellationToken) -> Fut, ...
{
  // odnawia wszystkie leasing_duration_secs / 3
  // z wyłączeniem K8: praca (token) tylko raz, token nigdy nie anulowany
}
#
Prawdziwy fan-out

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-frontCo 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
channels/manager.rs → nats/mod.rs (extraits réels)rust
// cdc-worker publikuje (ChannelManager::publish_and_relay)
nc.publish_with_headers(subject, headers, payload).await

// ws-front nasłuchuje całego poddrzewa (NatsRelay::start_relay)
client.subscribe("aura.realtime.>").await

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.

#
Bezpieczeństwo

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.

subscriptions.rs::check_rls_visible_batch (prawdziwy, uproszczony)sql
BEGIN;
SET LOCAL ROLE "aura_authenticated";
-- Roszczenia JWT abonenta wprowadzone do GUC
SELECT set_config('request.jwt.claims', $1, true);
SELECT EXISTS(SELECT 1 FROM "schema"."table" WHERE id::text = $2);
ROLLBACK;

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.

Astus

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.

#
Pułapka zweryfikowana

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.

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

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.

#
Aktualny limit

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.

#
Strona klienta

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.

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

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.

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

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.

GOTOWY DO WDROŻENIA?

Twój backend w pięć minut.

Karta kredytowa nie jest wymagana · 500 MB za darmo · 50 000 MAU