PRODसंप्रभु यूरोपीय BaaS मंचडैशबोर्ड खोलें →

इंजीनियरिंग · 10 मिनट पढ़ा

PostgreSQL CDC NATS के साथ: वास्तविक समय वास्तुकला

Affane Daylami · Fondateur · 24 जुलाई 2026

ब्लॉग पर वापस जाएँ

WebSocket के माध्यम से कनेक्टेड Aurabase क्लाइंट को COMMIT के कुछ मिलीसेकंड बाद बेस में डाली गई एक लाइन दिखाई देती है - बिना पोलिंग के, बिना वेबहुक को कॉन्फ़िगर किए। तंत्र तीन चरणों में है. एक तार्किक प्रतिकृति प्लगइन पोस्टग्रेज वाल को JSON में डिकोड करता है, केवल एक प्रक्रिया को इसे पढ़ने का अधिकार है, और NATS प्रत्येक घटना को क्लस्टर में सभी WebSocket सर्वर पर रिले करता है। यहां बताया गया है कि इसे ऑरा-रियलटाइम सर्विस के वास्तविक कोड में कैसे जोड़ा गया है।

यह अंग्रेजी पाठ फ़्रेंच मूल से स्वचालित रूप से उत्पन्न हुआ था और अभी तक इसकी समीक्षा नहीं की गई है।
यह पृष्ठ स्वचालित रूप से अनुवादित किया गया था. अंग्रेजी संस्करण प्रामाणिक है.

यहां कुछ भी आदर्शीकृत पैटर्न नहीं है। प्रत्येक उद्धरण प्रकाशन की तिथि पर aura-realtimeक्रेट से आता है। यहां एक स्पष्ट अंतर है: NATS जेटस्ट्रीम वह नहीं है जो इस पाइपलाइन में सीडीसी घटनाओं को ट्रांसपोर्ट करता है। हम नीचे विस्तार से बताते हैं कि जेटस्ट्रीम वास्तव में यहां क्या करता है, और इसके स्थान पर सीडीसी ट्रैफ़िक क्या होता है। एक सत्यापित कुबेरनेट्स ट्रैप, जो बिना किसी त्रुटि के पूरे तंत्र को शांत करने में सक्षम है, को नीचे प्रलेखित किया गया है।

अनिवार्य है
  • ऑराबेस का पोस्टग्रेज सीडीसी वाल को wal2json और pg_logical_slot_get_changes के माध्यम से पढ़ता है - pgoutputबाइनरी प्रोटोकॉल नहीं, डेबेज़ियम/काफ्का कनेक्ट ब्रिज नहीं।
  • तार्किक स्लॉट aura_cdc_slot केवल एक पाठक को सहन करता है: aura-realtime-cdc-worker को कुबेरनेट्स लीज (coordination.k8s.io/v1) के माध्यम से चुना जाता है, K8s को छोड़कर "हमेशा लीडर" मोड के साथ।
  • ws-front प्रतिकृतियों का फैन-आउट NATS कोर (aura.realtime.>पर सरल पब/उप) के माध्यम से होता है, न कि लगातार जेटस्ट्रीम स्ट्रीम के माध्यम से। जेटस्ट्रीम, इसी सेवा में, विशेष रूप से क्रॉस-इंस्टेंस उपस्थिति केवी की सेवा प्रदान करता है।
  • प्रत्येक ईवेंट को ट्रांसमिशन से ठीक पहले, संबंधित लाइन पर वास्तविक SET LOCAL ROLE अनुरोध के माध्यम से ग्राहक द्वारा आरएलएस में दोबारा जांचा जाता है - कैश्ड सन्निकटन नहीं।
  • रिपॉजिटरी इतिहास में जांच की गई गड़बड़ी: पोस्टग्रेज कुबेरनेट्स छवि में wal2json प्लगइन के बिना, स्लॉट निर्माण विफल हो जाता है और सीडीसी चुपचाप निष्क्रिय रहता है।
#
इंजन का चुनाव

wal2json, pgoutput नहीं, Debezium नहीं

अधिकांश सीडीसी पोस्टग्रेज पाइपलाइनें pgoutput, बाइनरी लॉजिक प्रतिकृति प्रोटोकॉल से गुजरती हैं, और फिर डेबेज़ियम जैसे कनेक्टर के माध्यम से जाती हैं जो इसे काफ्का में अनुवादित करता है। ऑराबेस इस चरण को छोड़ देता है: aura-realtime सेवा सीधे [WAL को wal2json प्लगइनके साथ तार्किक प्रतिकृति](https://www.postgresql.org/docs/16/logicaldecoding-explanation.html) में डिकोड करती है, जो मध्यवर्ती अनुवाद चरण के बिना प्रयोग करने योग्य JSON का उत्पादन करती है।

cdc/postgres.rs (वास्तविक प्रश्न)sql
-- अनुपस्थित होने पर एक बार बनाया गया, यदि वाल समाप्त हो गया तो पुनः बनाया गया
SELECT pg_create_logical_replication_slot(
  'aura_cdc_slot', 'wal2json');

-- बार-बार मतदान (बैच खाली आने पर 100 एमएस का बैकऑफ़)
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');

कोड में सत्यापित तीन पोस्टग्रेज़ पूर्वापेक्षाएँ: एप्लिकेशन भूमिका में REPLICATIONहोना चाहिए, wal_level = logical सक्रिय होना चाहिए, और max_slot_wal_keep_size में वाल प्रतिधारण को सीमित होना चाहिए। इस अंतिम सीमा के बिना, एक धीमा उपभोक्ता डिस्क की मात्रा को अनिश्चित काल तक बढ़ने का कारण बनता है। सेवा स्लॉट के wal_status पर भी नज़र रखती है: यदि यह lost में बदल जाता है (WAL इस सीमा से परे शुद्ध हो जाता है), तो स्लॉट स्वचालित रूप से पुनः निर्मित हो जाता है। इस बीच उपभोग नहीं की गई घटनाओं के नुकसान को मान लिया गया है और लॉग किया गया है।

प्रति मतदान 1000 परिवर्तनों का एक बैच, हर 3 सेकंड में एक गतिशील तालिका फ़िल्टर ताज़ा होता है। केवल वे तालिकाएँ जहाँ किसी प्रोजेक्ट ने वास्तविक समय को स्पष्ट रूप से सक्षम किया है, add-tables दर्ज करें - यह उन तालिकाओं के वाल को डिकोड करने से बचाता है जो किसी भी ग्राहक के लिए कोई रुचि नहीं रखती हैं।

#
एक एकल निर्माता

सीडीसी-कार्यकर्ता: कुबेरनेट्स लीज द्वारा चुना गया, डुप्लिकेट नहीं

एक पोस्टग्रेज़ लॉजिकल प्रतिकृति स्लॉट एक समय में केवल एक सक्रिय ड्राइव को सहन करता है। समान aura_cdc_slot के एकाधिक उपभोक्ताओं को समानांतर में चलाने से परिवर्तनों का क्रम टूट जाएगा, न कि केवल उनकी नकल होगी। वेबसॉकेट के क्षैतिज पैमाने का त्याग किए बिना इस बाधा को हल करने के लिए ऑराबेस ने aura-realtime को दो अलग-अलग बायनेरिज़ में विभाजित किया।

Cargo.tomltoml
# सीडीसी-कार्यकर्ता: निर्वाचित नेता, सीडीसी लॉन्च करता है + NATS पर प्रकाशित करता है। कोई WS सर्वर नहीं.
[[bin]]
name = "aura-realtime-cdc-worker"

# डब्ल्यूएस-फ्रंट: एन रेप्लिकास (एचपीए), एनएटीएस, आरएलएस फिल्टर का उपभोग करता है, डब्ल्यूएस/एसएसई परोसता है।
[[bin]]
name = "aura-realtime-ws-front"

cdc-worker केवल लीज कुबेरनेट्स ऑब्जेक्ट (coordination.k8s.io/v1) को पकड़कर स्लॉट को पढ़ने का अधिकार प्राप्त करता है। वह इसे बनाने का प्रयास करता है, या यदि यह समाप्त हो गया है तो इसे चुरा लेता है, फिर इसके जीवनकाल के एक तिहाई के बाद इसे नवीनीकृत करता है। जब पट्टा खो जाता है (नवीनीकरण विफल हो जाता है, धारक बदल जाता है), एक CancellationToken प्रगति में चल रहे कार्य को काट देता है और प्रक्रिया एक अधिग्रहण लूप पर वापस आ जाती है। कुबेरनेट्स के बाहर - स्थानीय परीक्षण या गैर-क्लस्टर विकास के दौरान, उदाहरण के लिए - K8s क्लाइंट कनेक्ट करने में विफल रहता है, और सेवा "हमेशा अग्रणी" मोड पर स्विच हो जाती है। विकास में उपयोगी, गलत है यदि आप इसे कई प्रतिकृतियों के साथ उत्पाद में भूल जाते हैं।

cdc/leader.rs (असली हस्ताक्षर)rust
pub async fn run_with_lease<F, Fut>(cfg: LeaseConfig, work: F) -> anyhow::Result<()>
where F: Fn(CancellationToken) -> Fut, ...
{
  // सभी लीज़_अवधि_सेकंड/3 को नवीनीकृत करता है
  // K8s को छोड़कर: कार्य (टोकन) केवल एक बार, टोकन कभी रद्द नहीं किया जाता
}
#
असली फैन-आउट

सीडीसी के लिए NATS कोर, उपस्थिति के लिए जेटस्ट्रीम

इस प्रकार की पाइपलाइन में यह सबसे अधिक गलत समझा जाने वाला बिंदु है: NATS और JetStream दो अलग-अलग चीज़ें हैं, और CDC ईवेंट लगातार JetStream स्ट्रीम से नहीं गुजरते हैं। cdc-worker प्रत्येक ईवेंट को Client::publish_with_headers - NATS कोर पब/सब एपीआईके साथ प्रकाशित करता है, जो जेटस्ट्रीम से अलग, बिना किसी दृढ़ता या रीप्ले के को अधिकतम एक बार डिलीवर करता है।

NATS हार्ट (पब/उप)डब्ल्यूएस-फ्रंट प्रतिकृतियों के लिए सीडीसी घटनाओं का फैन-आउटअधिक से अधिक एक बार, बिना किसी दृढ़ता या दोहराव के
NATS जेटस्ट्रीम (KV)क्रॉस-इंस्टेंस उपस्थिति (aura_presence)साझा स्थिति 60, दोबारा चलाने के लिए सीडीसी स्ट्रीम नहीं
channels/manager.rs → nats/mod.rs (extraits réels)rust
// सीडीसी-कार्यकर्ता प्रकाशित करता है (ChannelManager::publish_and_relay)
nc.publish_with_headers(subject, headers, payload).await

// डब्ल्यूएस-फ्रंट संपूर्ण सबट्री को सुनता है (NatsRelay::start_relay)
client.subscribe("aura.realtime.>").await

प्रत्येक प्रतिकृति ws-front सबट्री aura.realtime.>की सदस्यता लेती है, विषय को मूल चैनल नाम पर डीकोड करती है। इसके बाद यह इवेंट को अपने स्थानीय tokio::sync::broadcastपर, इससे जुड़े वेबसॉकेट के लिए पुनः प्रकाशित करता है। एक Aura-Origin हेडर भेजने वाले उदाहरण के पहचानकर्ता को रखता है: प्रत्येक प्रतिकृति उन संदेशों को अनदेखा करती है जिन्हें उसने स्वयं प्रकाशित किया है, जो केंद्रीय समन्वय के बिना डुप्लिकेट से बचाता है।

इस विकल्प का सीधा परिणाम होता है: NATS कोर लंबे समय तक चलने वाली कतार के बिना, सबसे अच्छा एक बार (at-most-one__) वितरित करता है। यदि किसी घटना के गुजरने पर ws-front प्रतिकृति थोड़ी देर के लिए NATS से डिस्कनेक्ट हो जाती है, तो यह पकड़ में नहीं आती है। प्रति-चैनल इतिहास (Channel.history, HISTORY_SIZEसे घिरा एक इन-मेमोरी बफर) प्रत्येक प्रतिकृति पर स्थानीय रूप से रहता है, क्लस्टर स्तर पर नहीं। यह एक स्वीकार्य समझौता है क्योंकि पोस्टग्रेज तार्किक प्रतिकृति स्लॉट सत्य का अपस्ट्रीम टिकाऊ स्रोत बना हुआ है। यह wal2json है जो यह सुनिश्चित करता है कि उपभोग से पहले कोई DB परिवर्तन नष्ट न हो, NATS नहीं।

जेटस्ट्रीम aura-realtime में मौजूद है - लेकिन पूरी तरह से अलग उपयोग के लिए। यह क्रॉस-इंस्टेंस उपस्थिति केवी स्टोर (aura_presence, मेमोरी स्टोर, 60-सेकंड max_age) पर कार्य करता है, जो सभी ws-frontप्रतिकृतियों के बीच किस चैनल पर ऑनलाइन है, इसे सिंक्रनाइज़ करता है। एक अल्पकालिक साझा स्थिति, दोबारा चलाने के लिए सीडीसी घटनाओं की एक धारा नहीं।

#
सुरक्षा

प्रत्येक इवेंट में ग्राहक द्वारा आरएलएस की दोबारा जांच की जाती है

एक सीडीसी कार्यक्रम किसी चैनल के सभी ग्राहकों के लिए अप्रमाणित नहीं होता है। पोलर रूट प्रत्येक को CDC_NUM_SHARDS कार्यकर्ताओं में से एक में बदल देता है (डिफ़ॉल्ट रूप से 4, project_idद्वारा हैश किया गया), जो पहले pg_policiesको पोल करता है। आरएलएस नीति के बिना एक तालिका सभी ग्राहकों को - सुपाबेस व्यवहार - की अनुमति देती है, जबकि नीतियों वाली एक तालिका प्रति-ग्राहक जांच को ट्रिगर करती है।

subscriptions.rs::check_rls_visible_batch (वास्तविक, सरलीकृत)sql
BEGIN;
SET LOCAL ROLE "aura_authenticated";
-- ग्राहक के JWT दावे GUC में डाले गए
SELECT set_config('request.jwt.claims', $1, true);
SELECT EXISTS(SELECT 1 FROM "schema"."table" WHERE id::text = $2);
ROLLBACK;

लेन-देन ग्राहक की भौतिक पोस्टग्रेज़ भूमिका को अपनाता है और उसके JWT दावों को request.jwt.claims में इंजेक्ट करता है - जिसे auth.uid() और auth.role() पॉलिसी पक्ष पर पढ़ते हैं। फिर यह जाँचता है कि रेखा इस भूमिका के अंतर्गत दृश्यमान रहती है, फिर रद्द कर देती है: कोई लेखन नहीं, एक वास्तविक आरएलएस पढ़ा गया। घटना के allowed_sub_ids में केवल मान्य sub_id भूमि; खाली वीईसी का मतलब फेल-क्लोज्ड में कोई नहीं है।

अस्टुसे

एक DELETE प्रति पंक्ति इस जांच से बच जाता है - लाइन गायब हो गई है, इसे आरएलएस के तहत दोबारा पढ़ना असंभव है - और चैनल के सभी ग्राहकों के पास वापस चला जाता है। बदले में, DELETE का पेलोड केवल पहचान कॉलम (प्राथमिक कुंजी) को उजागर करता है, कभी भी हटाई गई पंक्ति की सामग्री को नहीं।

एक विवरण जो अक्सर माइग्रेशन के रास्ते में आता है: wal2json केवल UPDATE या DELETE की पूरी सामग्री को कैप्चर करता है यदि तालिका में REPLICA IDENTITY FULLहै। इसके बिना, केवल INSERT वापस आते हैं। यही कारण है कि तालिका पर वास्तविक समय को सक्षम करने वाला समापन बिंदु इस ALTER TABLE को स्वचालित रूप से सेट करता है - रिपॉजिटरी के लिए एक वास्तविक फिक्स, एक अलग चेकबॉक्स नहीं।

#
जाल का सत्यापन किया गया

कुबेरनेट्स में सीडीसी चुप क्यों रह सकता है?

फाइलिंग का इतिहास एक वास्तविक जाल का दस्तावेजीकरण करता है, सैद्धांतिक मामला नहीं। कुबेरनेट्स क्लस्टर (pgvector/pgvector स्टॉक समुदाय छवि) पर डिफ़ॉल्ट रूप से उपयोग की जाने वाली पोस्टग्रेज छवि में wal2jsonप्लगइन शामिल नहीं है। इसके बिना, pg_create_logical_replication_slot('aura_cdc_slot', 'wal2json') विफल हो जाता है, और क्लाइंट पक्ष पर कुछ भी स्पष्ट रूप से इसका संकेत नहीं देता है। वेबसॉकेट खुले रहते हैं, सदस्यताएँ स्वीकार की जाती हैं, लेकिन कोई DB ईवेंट कभी नहीं होता है।

docker/Postgres.Dockerfile (असली अर्क)dockerfile
RUN apt-get update && \
    apt-get install -y --no-install-recommends \
      postgresql-16-wal2json=2.6-4.pgdg12+1 \
      ...

वास्तविक समाधान इस पैकेज को स्थापित करके एक समर्पित पोस्टग्रेज छवि बनाना और प्रकाशित करना था, फिर स्टॉक छवि के बजाय कुबेरनेट्स StatefulSet को उस पर इंगित करें। यह स्पष्ट रूप से wal_level=logical, max_replication_slots और max_slot_wal_keep_size को स्टार्टअप तर्क के रूप में सेट करता है - सामान्य छवियों से अनुपस्थित।

इस प्रकार की विफलता को देखने योग्य बनाने के लिए ही cdc-worker अपने समर्पित /health समापन बिंदु पर एक cdc_active ध्वज (जब तक स्लॉट स्वस्थ होने की पुष्टि नहीं की जाती है) को उजागर करता है। यह देखने के लिए एक द्विआधारी संकेत है - केवल क्लाइंट-साइड घटनाओं की अनुपस्थिति से मृत सीडीसी का अनुमान लगाने के बजाय।

#
वर्तमान सीमा

यह पाइपलाइन अभी तक क्या कवर नहीं करती है: समर्पित क्लस्टर

यह तंत्र एक एकल POSTGRES_REPLICATION_URLवैरिएबल को पढ़ता है, इसलिए एक एकल पोस्टग्रेज सर्वर और एक एकल लॉजिकल स्लॉट। यह साझा क्लस्टर पर प्रति प्रोजेक्ट मॉडल (project_<uuid>) मल्टी-टेनेंट स्कीमा के अनुरूप है। इस मॉडल पर, सब कुछ काम करता है: एक एकल cdc-worker एक ही क्लस्टर में सभी परियोजनाओं के परिवर्तनों को देखता है, और आरेख द्वारा रूट करता है।

एक समर्पित सीएनपीजी क्लस्टर पर एक परियोजना के लिए - इसका अपना पृथक पोस्टग्रेज उदाहरण - आज उपयोग की गई छवि (क्लाउडनेटिवपीजी स्टॉक छवि से निर्मित) केवल pg_graphqlजोड़ती है, wal2jsonनहीं। कुछ भी नहीं, प्रति समर्पित क्लस्टर cdc-worker चलाता है। कोड की वर्तमान स्थिति में, यहां वर्णित वास्तविक समय सीडीसी साझा स्तर पर केंद्रित है - एक वास्तविक वास्तुशिल्प सीमा, न कि रोडमैप कार्यक्षमता।

#
ग्राहक पक्ष

क्या नहीं बदलता: SDK का postgres_changes API

यह संपूर्ण तंत्र SDK से अदृश्य रहता है।aurabase-js चेनेबल एपीआई इस पर निर्भर नहीं करता है कि घटना NATS के माध्यम से cdc-worker से आती है या, स्थानीय डेव में, अनस्प्लिट aura-realtime बाइनरी से आती है।

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

लेकिन एक भी घटना घटित होने से पहले, तालिका को सीडीसी के साथ पंजीकृत होना चाहिए - यह निर्माण में स्वचालित नहीं है। समर्पित एंडपॉइंट (या स्टूडियो में समतुल्य टॉगल) पर एक service_role कॉल प्रोजेक्ट के प्लेटफ़ॉर्म स्कीमा के realtime_tables में तालिका लिखता है, और आपके लिए 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}'

बाकी सुपाबेस संगत एसडीके सतह - प्रमाणीकरण, भंडारण, नीतियां आरएलएस - के लिए, माइग्रेशन गाइडदेखें। aura-realtime एक एकल कार्यस्थान क्रेट में एकाधिक बायनेरिज़ में कैसे विभाजित होता है, इसके लिए कार्गो आर्किटेक्चर विवरणदेखें।

तैनाती के लिए तैयार हैं?

पाँच मिनट में आपका बैकएंड।

किसी क्रेडिट कार्ड की आवश्यकता नहीं · 500 एमबी निःशुल्क · 50,000 एमएयू