PRODEgemen Avrupa BaaS platformuKontrol Panelini Aç →

Mühendislik · 10 dk. okuma

PostgreSQL CDC ile NATS: gerçek zamanlı mimari

Affane Daylami · Fondateur · 24 Temmuz 2026

Bloga geri dön

WebSocket aracılığıyla bağlanan bir Aurabase istemcisi, COMMIT'den birkaç milisaniye sonra tabana eklenen bir satırın göründüğünü görüyor; yoklama olmadan, yapılandırılacak web kancası olmadan. Mekanizma üç aşamalıdır. Mantıksal bir çoğaltma eklentisi, Postgres WAL'ın kodunu JSON'a çözer, yalnızca bir işlemin onu okuma hakkı vardır ve NATS, her olayı kümedeki tüm WebSocket sunucularına aktarır. Aura-realtimeservice'in gerçek kodunda bu şekilde bağlanmıştır.

Bu İngilizce metin, Fransızca orijinalinden otomatik olarak oluşturulmuştur ve henüz incelenmemiştir.
Bu sayfa otomatik olarak çevrildi. İngilizce versiyonu yetkilidir.

Buradaki hiçbir şey idealize edilmiş bir model değildir. Her alıntı yayın tarihinde aura-realtimekutusundan olduğu gibi gelir. Burada net bir ayrım var: NATS JetStream bu kanalda CDC olaylarını aktaran şey değil. JetStream'in burada gerçekte ne yaptığını ve CDC trafiğini yerinde neyin taşıdığını aşağıda ayrıntılarıyla anlatıyoruz. Herhangi bir hataya yol açmadan tüm mekanizmayı sessiz hale getirebilen doğrulanmış bir Kubernetes tuzağı aşağıda daha ayrıntılı olarak belgelenmiştir.

Temeller
  • Aurabase'in Postgres CDC'si WAL'yi wal2json ve pg_logical_slot_get_changes aracılığıyla okur; pgoutputikili protokolü veya Debezium/Kafka Connect köprüsü aracılığıyla değil.
  • aura_cdc_slot mantıksal yuvası yalnızca bir okuyucuyu tolere eder: aura-realtime-cdc-worker, K8'leri hariç tutan "her zaman lider" moduyla bir Kubernetes Kiralaması (coordination.k8s.io/v1) aracılığıyla seçilir.
  • ws-front kopyalarına dağıtım, kalıcı bir JetStream akışı aracılığıyla değil, NATS çekirdeği (aura.realtime.>üzerinde basit pub/sub) aracılığıyla yapılır. JetStream, aynı hizmette yalnızca örnekler arası varlık KV'sine hizmet eder.
  • Her olay, iletimden hemen önce, önbelleğe alınmış bir yaklaşım değil, ilgili hattaki gerçek bir SET LOCAL ROLE isteği yoluyla abone tarafından RLS'de yeniden kontrol edilir.
  • Depo geçmişinde kontrol edilen tehlike: Postgres Kubernetes görüntüsünde wal2json eklentisi olmadan yuva oluşturma işlemi başarısız olur ve CDC sessizce devre dışı kalır.
#
Motor seçimi

wal2json, pgoutput değil, Debezium değil

Çoğu CDC Postgres işlem hattı, ikili mantık çoğaltma protokolü olan pgoutputüzerinden ve ardından onu Kafka'ya çeviren Debezium gibi bir bağlayıcı üzerinden geçer. Aurabase bu adımı atlar: aura-realtime hizmeti, WAL'nin kodunu doğrudan ile wal2json eklentisiile mantıksal çoğaltmaya dönüştürür ve ara çeviri aşaması olmadan kullanılabilir JSON üretir.

cdc/postgres.rs (gerçek sorgular)sql
-- Yoksa bir kez oluşturulur, WAL'ın süresi dolarsa yeniden oluşturulur
SELECT pg_create_logical_replication_slot(
  'aura_cdc_slot', 'wal2json');

-- Tekrarlanan anket (toplu iş boş döndüğünde 100 ms geri çekilme)
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');

Kodda doğrulanan üç Postgres önkoşulu: uygulama rolü REPLICATIONtaşımalı, wal_level = logical etkin olmalı ve max_slot_wal_keep_size WAL saklamayı sınırlamalıdır. Bu son sınır olmadan, yavaş bir tüketici disk hacminin süresiz olarak büyümesine neden olur. Hizmet ayrıca yuvanın wal_status öğesini de izler: lost (WAL bu sınırın ötesinde temizlenir) olarak değişirse yuva otomatik olarak yeniden oluşturulur. Bu arada tüketilmeyen olayların kaybı varsayılır ve günlüğe kaydedilir.

Her 3 saniyede bir yenilenen dinamik tablo filtresiyle anket başına 1000 değişiklikten oluşan bir grup. Yalnızca bir projenin açıkça gerçek zamanı etkinleştirdiği tablolar add-tables değerini girer; bu, herhangi bir abonenin ilgisini çekmeyen tabloların WAL kodunun çözülmesini önler.

#
Tek bir yapımcı

cdc-worker: Kubernetes Lease tarafından seçilir, kopyalanmaz

Postgres mantıksal çoğaltma yuvası aynı anda yalnızca bir etkin sürücüyü tolere eder. Aynı aura_cdc_slot öğesinin birden fazla tüketicisini paralel olarak çalıştırmak, yalnızca bunları çoğaltmak değil, değişikliklerin sırasını da bozar. Aurabase, WebSocket'in yatay ölçeğinden ödün vermeden bu kısıtlamayı çözmek için aura-realtime'yi iki ayrı ikili dosyaya böldü.

Cargo.tomltoml
# cdc-worker: seçilmiş lider, CDC+'yi başlatır ve NATS'te yayınlar. WS sunucusu yok.
[[bin]]
name = "aura-realtime-cdc-worker"

# ws-front: N kopya (HPA), NATS, RLS filtresi tüketir, WS/SSE'ye hizmet eder.
[[bin]]
name = "aura-realtime-ws-front"

cdc-worker yalnızca bir Lease Kubernetes nesnesi (coordination.k8s.io/v1) tutarak yuvayı okuma hakkını kazanır. Onu yaratmaya çalışır veya süresi dolmuşsa çalmaya çalışır, sonra ömrünün üçte biri sonunda onu yeniler. Kira sözleşmesi kaybolduğunda (yenileme başarısız oldu, sahibi değişti), CancellationToken devam eden işi keser ve süreç bir satın alma döngüsüne geri döner. Kubernetes dışında (örneğin, yerel test veya küme dışı geliştirme sırasında) K8s istemcisi bağlanamıyor ve hizmet "her zaman önde" moduna geçiyor. Geliştirmede kullanışlıdır, birkaç kopyayla üretimde unutursanız yanlış olur.

cdc/leader.rs (gerçek imzalar)rust
pub async fn run_with_lease<F, Fut>(cfg: LeaseConfig, work: F) -> anyhow::Result<()>
where F: Fn(CancellationToken) -> Fut, ...
{
  // tüm rent_duration_secs / 3'ü yeniler
  // K8'ler hariç: yalnızca bir kez çalışır (belirteç), belirteç hiçbir zaman iptal edilmez
}
#
Gerçek yayılma

CDC için NATS çekirdeği, varlık için JetStream

Bu, bu tür bir işlem hattında en sık yanlış anlaşılan noktadır: NATS ve JetStream iki ayrı şeydir ve CDC olayları kalıcı bir JetStream akışından geçmez. cdc-worker her olayı Client::publish_with_headers — NATS çekirdek yayın/alt API'siile yayınlar; bu, JetStream'den ayrı olarak, kalıcılık veya tekrar oynatma olmadan en fazla bir kez teslim edilir.

NATS kalbi (pub/sub)CDC olaylarının ws-ön kopyalarına yayılmasıEn fazla bir kez, ısrar etmeden veya tekrarlama olmadan
NATS JetStream (KV)Örnekler arası mevcudiyet (aura_presence)Paylaşılan durum 60'lar, tekrar oynatılacak bir CDC akışı değil
channels/manager.rs → nats/mod.rs (extraits réels)rust
// cdc-worker yayınlıyor (ChannelManager::publish_and_relay)
nc.publish_with_headers(subject, headers, payload).await

// ws-front alt ağacın tamamını dinler (NatsRelay::start_relay)
client.subscribe("aura.realtime.>").await

Her kopya ws-front aura.realtime.>alt ağacına abone olur ve konunun kodunu orijinal kanal adına çözer. Daha sonra olayı, kendisine bağlı WebSocket'ler için yerel tokio::sync::broadcastkonumunda yeniden yayınlar. Bir Aura-Origin başlığı, gönderen örneğin tanımlayıcısını taşır: her kopya, kendisinin yayınladığı mesajları yok sayar; bu, merkezi koordinasyon olmaksızın kopyaların önlenmesini sağlar.

Bu seçimin doğrudan bir sonucu vardır: NATS çekirdeği uzun süreli bir kuyruk olmadan en iyi şekilde bir kez (en fazla bir kez) dağıtım yapar. Bir olay geçtiğinde bir ws-front kopyasının NATS ile bağlantısı kısa süreliğine kesilirse, olayı yakalamaz. Kanal başına geçmiş (Channel.history, HISTORY_SIZEile sınırlanan bir bellek içi arabellek), küme düzeyinde değil, her çoğaltmada yerel olarak bulunur. Bu kesinlikle kabul edilebilir bir uzlaşmadır çünkü Postgres mantıksal çoğaltma yuvası, yukarı yöndeki dayanıklı gerçeğin kaynağı olmaya devam etmektedir. NATS değil, tüketimden önce hiçbir veritabanı değişikliğinin kaybolmamasını sağlayan wal2json'dir.

JetStream aura-realtime'de mevcut, ancak tamamen farklı bir kullanım için. Tüm ws-frontkopyaları arasında kimin hangi kanalda çevrimiçi olduğunu senkronize eden çapraz örnek varlığı KV deposuna (aura_presence, bellek deposu, 60 saniyelik max_age) hizmet eder. Kısa ömürlü bir paylaşılan durum, tekrarlanacak CDC olaylarının akışı değil.

#
Güvenlik

RLS her etkinlikte abone tarafından yeniden kontrol edilir

Bir CDC etkinliği, bir kanalın tüm abonelerine doğrudan ulaşmaz. Yoklayıcı, her değişikliği CDC_NUM_SHARDS çalışanlarından birine yönlendirir (varsayılan olarak 4, project_idile hashlenir), ilk yoklayan pg_policies. RLS politikası olmayan bir tablo tüm abonelere (Supabase davranışı) izin verirken, politikaları olan bir tablo abone başına kontrolü tetikler.

subscriptions.rs::check_rls_visible_batch (gerçek, basitleştirilmiş)sql
BEGIN;
SET LOCAL ROLE "aura_authenticated";
-- GUC'ye enjekte edilen abonenin JWT talepleri
SELECT set_config('request.jwt.claims', $1, true);
SELECT EXISTS(SELECT 1 FROM "schema"."table" WHERE id::text = $2);
ROLLBACK;

İşlem, abonenin fiziksel Postgres rolünü üstlenir ve JWT taleplerini, politika tarafında auth.uid() ve auth.role() tarafından okunan request.jwt.claims içine enjekte eder. Daha sonra satırın bu rol altında görünür durumda kalıp kalmadığını kontrol eder ve ardından iptal eder: yazma yok, gerçek bir RLS okuması. Yalnızca doğrulanmış sub_id etkinliğin allowed_sub_ids alanına girer; boş bir vec, arıza kapalı durumunda hiç kimse anlamına gelmez.

Astuce

Bir DELETE satır başına bu kontrolden kaçar - satır kaybolur, onu RLS altında yeniden okumak imkansızdır - ve kanalın tüm abonelerine geri döner. Buna karşılık, bir DELETE'in yükü yalnızca kimlik sütunlarını (birincil anahtar) açığa çıkarır, silinen satırın içeriğini asla açığa çıkarmaz.

Geçişin önünde sıklıkla engel olan bir ayrıntı: wal2json, yalnızca tabloda REPLICA IDENTITY FULLvarsa UPDATE veya DELETE öğesinin tüm içeriğini yakalar. Bu olmadan yalnızca INSERT geri gelir. Bu nedenle, bir tabloda gerçek zamanlıyı etkinleştiren uç nokta, bu ALTER TABLE öğesini otomatik olarak ayarlar; bu, ayrı bir onay kutusu değil, depoya yönelik gerçek bir düzeltmedir.

#
Tuzak doğrulandı

CDC neden Kubernetes'te sessiz kalabiliyor?

Dosyalamanın geçmişi teorik bir vakayı değil, gerçek bir tuzağı belgeliyor. Bir Kubernetes kümesinde varsayılan olarak kullanılan Postgres görüntüsü (pgvector/pgvector hazır topluluk görüntüsü), wal2jsoneklentisini içermez. Bu olmadan, pg_create_logical_replication_slot('aura_cdc_slot', 'wal2json') başarısız olur ve istemci tarafındaki hiçbir şey bunu açıkça belirtmez. WebSocket'ler açık kalır, abonelikler kabul edilir, ancak hiçbir DB olayı gerçekleşmez.

docker/Postgres.Dockerfile (gerçek özü)dockerfile
RUN apt-get update && \
    apt-get install -y --no-install-recommends \
      postgresql-16-wal2json=2.6-4.pgdg12+1 \
      ...

Asıl düzeltme, bu paket yüklüyken özel bir Postgres görüntüsü oluşturmak ve yayınlamak, ardından hazır görüntü yerine Kubernetes StatefulSet'yi ona yönlendirmekti. Ayrıca, genel görüntülerde bulunmayan wal_level=logical, max_replication_slots ve max_slot_wal_keep_size öğelerini başlangıç ​​bağımsız değişkenleri olarak açıkça ayarladı.

cdc-worker'nin özel /health uç noktasında bir cdc_active bayrağını (yuvanın sağlıklı olduğu doğrulanmadığı sürece yanlış) göstermesi tam olarak bu tür bir arızayı gözlemlenebilir kılmak içindir. Bu, yalnızca istemci tarafı olayların yokluğundan ölü CDC'yi çıkarmak yerine, izlenmesi gereken ikili bir sinyaldir.

#
Akım sınırı

Bu hattın henüz kapsamadığı konular: özel kümeler

Bu mekanizma tek bir POSTGRES_REPLICATION_URLdeğişkenini, dolayısıyla tek bir Postgres sunucusunu ve tek bir mantıksal yuvayı okur. Bu, paylaşılan bir kümedeki proje modeli başına çok kiracılı şema (project_<uuid>) ile tutarlıdır. Bu modelde her şey çalışır: tek bir cdc-worker aynı kümedeki tüm projelerdeki değişiklikleri diyagram bazında görür.

Özel bir CNPG kümesindeki (kendi yalıtılmış Postgres örneği) bir proje için, bugün kullanılan görüntü (CloudNativePG stok görüntüsünden oluşturulmuştur) wal2jsondeğil, yalnızca pg_graphqlekler. Hiçbir şey de ayrılmış küme başına cdc-worker çalıştırmaz. Kodun mevcut durumunda, burada açıklanan gerçek zamanlı CDC, bu nedenle, bir yol haritası işlevi değil, gerçek bir mimari sınır olan paylaşılan katmana odaklanmaya devam eder.

#
İstemci tarafı

Değişmeyen şey: SDK'nın postgres_changes API'si

Bu mekanizmanın tamamı SDK'dan görünmez kalır.aurabase-js zincirlenebilir API, olayın NATS aracılığıyla cdc-worker'den mi yoksa yerel geliştirmede bölünmemiş aura-realtime ikili dosyasından mı geldiğine bağlı olarak değişmez.

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

Ancak tek bir olayın gerçekleşebilmesi için tablonun CDC'ye kaydedilmesi gerekir; bu, oluşturma sırasında otomatik olarak gerçekleşmez. Özel uç noktadaki (veya Studio'daki eşdeğer geçiş) bir service_role çağrısı, projenin platform şemasının realtime_tables içindeki tabloyu yazar ve sizin için REPLICA IDENTITY FULL değerini ayarlar.

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

Supabase uyumlu SDK yüzeyinin geri kalanı (kimlik doğrulama, depolama, RLS politikaları) için geçiş kılavuzu'ye bakın. aura-realtime'nin tek bir çalışma alanı kasasında birden fazla ikili dosyaya nasıl bölündüğünü öğrenmek için Kargo mimarisi ayrıntısına bakın.

DAĞITILMAYA HAZIR MISINIZ?

Beş dakika içinde arka ucunuz.

Kredi kartı gerekmez · 500 MB ücretsiz · 50.000 MAU