PRODمنصة BaaS الأوروبية السياديةافتح لوحة المعلومات →

الهندسة · 10 دقيقة للقراءة

PostgreSQL CDC مع NATS: بنية الوقت الفعلي

Affane Daylami · Fondateur · 24 يوليو 2026

العودة إلى بلوق

يرى عميل Aurabase المتصل عبر WebSocket أن الخط المدرج في القاعدة يظهر بعد بضعة أجزاء من الثانية من COMMIT - بدون استقصاء، بدون خطاف ويب للتكوين. الآلية تتم على ثلاث مراحل. يقوم المكون الإضافي للنسخ المتماثل المنطقي بفك تشفير Postgres WAL إلى JSON، ويكون لعملية واحدة فقط الحق في قراءتها، وتقوم NATS بترحيل كل حدث إلى جميع خوادم WebSocket في المجموعة. وإليك كيفية ربطها بالرمز الفعلي لخدمة aura-realtimeservice.

تم إنشاء هذا النص الإنجليزي تلقائيًا من النص الأصلي الفرنسي ولم تتم مراجعته بعد.
تمت ترجمة هذه الصفحة تلقائيًا. النسخة الإنجليزية موثوقة.

لا شيء هنا هو نمط مثالي. يأتي كل مقتطف كما هو من صندوق aura-realtime، في تاريخ النشر. يوجد تمييز واضح هنا: NATS JetStream ليس هو ما ينقل أحداث CDC في مسار التدفق هذا. نوضح أدناه بالتفصيل ما يفعله JetStream بالفعل هنا، وما الذي يحمل حركة مرور CDC في مكانها. تم توثيق فخ Kubernetes الذي تم التحقق منه، والقادر على جعل الآلية بأكملها صامتة دون إثارة أي أخطاء، أدناه.

الأساسيات
  • يقوم Postgres CDC الخاص بـ Aurabase بقراءة WAL عبر wal2json وpg_logical_slot_get_changes - وليس البروتوكول الثنائي pgoutput، وليس جسر Debezium/Kafka Connect.
  • الفتحة المنطقية aura_cdc_slot تتسامح مع قارئ واحد فقط: aura-realtime-cdc-worker يتم اختياره عبر Kubernetes Lease (coordination.k8s.io/v1)، مع وضع "القائد دائمًا" باستثناء K8s.
  • يمر التوزيع الموسع للنسخ المتماثلة ws-front عبر NATS core (نشر/فرع بسيط على aura.realtime.>)، وليس من خلال دفق JetStream المستمر. JetStream، في هذه الخدمة نفسها، يخدم حصريًا التواجد عبر المثيلات KV.
  • يتم إعادة فحص كل حدث في RLS بواسطة المشترك، قبل الإرسال مباشرةً، عبر طلب SET LOCAL ROLE حقيقي على السطر المعني - وليس تقريبًا مخبأً.
  • تم التحقق من المأزق في سجل المستودع: بدون المكون الإضافي wal2json في صورة Postgres Kubernetes، يفشل إنشاء الفتحة ويظل مركز السيطرة على الأمراض (CDC) غير نشط بصمت.
#
اختيار المحرك

wal2json، وليس pgoutput، وليس Debezium

تمر معظم خطوط أنابيب CDC Postgres عبر pgoutput، وهو بروتوكول النسخ المتماثل المنطقي الثنائي، ثم عبر موصل مثل Debezium الذي يترجمه إلى Kafka. تتخطى Aurabase هذه الخطوة: تقوم خدمة aura-realtime بفك تشفير WAL مباشرة إلى النسخ المتماثل المنطقي مع المكون الإضافي wal2json، الذي ينتج JSON قابل للاستخدام دون مرحلة ترجمة وسيطة.

cdc/postgres.rs (استفسارات حقيقية)sql
-- يتم إنشاؤه مرة واحدة في حالة عدم وجوده، ويتم إعادة إنشائه في حالة انتهاء صلاحية WAL
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');

تم التحقق من ثلاثة متطلبات أساسية لـ Postgres في الكود: يجب أن يحمل دور التطبيق REPLICATION، ويجب أن يكون wal_level = logical نشطًا، وmax_slot_wal_keep_size يجب أن يحد من الاحتفاظ بـ WAL. بدون هذا الحد الأخير، يؤدي المستهلك البطيء إلى زيادة حجم القرص إلى أجل غير مسمى. تراقب الخدمة أيضًا wal_status للفتحة: إذا تغيرت إلى lost (تم إزالة WAL بعد هذا الحد)، فسيتم إعادة إنشاء الفتحة تلقائيًا. يتم افتراض وتسجيل فقدان الأحداث التي لم يتم استهلاكها في هذه الأثناء.

مجموعة مكونة من 1000 تغيير لكل استطلاع، مع تحديث مرشح الجدول الديناميكي كل 3 ثوانٍ. أدخل فقط الجداول التي قام المشروع بتمكين الوقت الفعلي بشكل صريح بإدخال add-tables - وهذا يتجنب فك تشفير WAL للجداول التي لا تهم أي مشترك.

#
منتج واحد

cdc-worker: يتم انتخابه بواسطة Kubernetes Lease، وليس مكررًا

تتسامح فتحة النسخ المتماثل المنطقية Postgres مع محرك أقراص نشط واحد فقط في كل مرة. سيؤدي تشغيل عدة مستهلكين لنفس aura_cdc_slot بالتوازي إلى كسر ترتيب التغييرات، وليس مجرد تكرارها. قامت Aurabase بتقسيم aura-realtime إلى ثنائيتين منفصلتين لحل هذا القيد دون التضحية بالمقياس الأفقي لـ WebSocket.

Cargo.tomltoml
# cdc-worker: زعيم منتخب، يطلق CDC+ وينشر على NATS. لا يوجد خادم WS.
[[bin]]
name = "aura-realtime-cdc-worker"

# ws-front: النسخ المتماثلة N (HPA)، تستهلك مرشح NATS، RLS، وتخدم WS/SSE.
[[bin]]
name = "aura-realtime-ws-front"

cdc-worker لا يكتسب سوى الحق في قراءة الفتحة من خلال الضغط على لتأجير كائن Kubernetes (coordination.k8s.io/v1). يحاول خلقها، أو سرقتها إذا انتهت صلاحيتها، ثم يجددها بعد ثلث عمرها. عند فقدان عقد الإيجار (فشل التجديد، تغيير المالك)، يقوم CancellationToken بقطع العمل الجاري وتعود العملية إلى حلقة الاستحواذ. خارج Kubernetes - أثناء الاختبار المحلي أو التطوير غير العنقودي، على سبيل المثال - يفشل عميل 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 الأساسية لمراكز السيطرة على الأمراض والوقاية منها (CDC)، وJetStream للتواجد

هذه هي النقطة التي يساء فهمها غالبًا في هذا النوع من خطوط الأنابيب: NATS وJetStream شيئان منفصلان، ولا تمر أحداث CDC عبر دفق JetStream المستمر. ينشر cdc-worker كل حدث باستخدام Client::publish_with_headers — واجهة NATS الأساسية/واجهة برمجة التطبيقات الفرعية، والذي يتم تسليمه مرة واحدة على الأكثر بدون استمرار أو إعادة تشغيل، منفصل عن JetStream.

قلب ناتس (حانة/فرعية)انتشار أحداث CDC إلى النسخ المتماثلة لـ ws-frontعلى الأكثر مرة واحدة، دون إصرار أو إعادة
ناتس جيت ستريم (كيه في)الحضور عبر المثيلات (aura_presence)الحالة المشتركة في الستينيات، وليس بث CDC لإعادة التشغيل
channels/manager.rs → nats/mod.rs (extraits réels)rust
// ينشر cdc-worker (ChannelManager::publish_and_relay)
nc.publish_with_headers(subject, headers, payload).await

// يستمع ws-front إلى الشجرة الفرعية بأكملها (NatsRelay::start_relay)
client.subscribe("aura.realtime.>").await

تشترك كل نسخة متماثلة ws-front في الشجرة الفرعية aura.realtime.>، وتقوم بفك تشفير الموضوع إلى اسم القناة الأصلية. ثم يقوم بإعادة نشر الحدث على موقعه المحلي tokio::sync::broadcastلـ WebSockets المتصلة به. يحمل رأس Aura-Origin معرف نسخة الإرسال: تتجاهل كل نسخة متماثلة الرسائل التي نشرتها بنفسها، مما يتجنب التكرارات بدون تنسيق مركزي.

هذا الاختيار له نتيجة مباشرة: يقوم مركز NATS بالتسليم مرة واحدة في أفضل الأحوال (مرة واحدة على الأكثر)، دون قائمة انتظار طويلة الأمد. إذا تم فصل نسخة متماثلة ws-front لفترة وجيزة عن NATS عند مرور حدث ما، فلن يتم اللحاق بها. التاريخ لكل قناة (Channel.history، مخزن مؤقت في الذاكرة يحده HISTORY_SIZE) موجود محليًا في كل نسخة متماثلة، وليس على مستوى المجموعة. يعد هذا حلاً وسطًا مقبولًا على وجه التحديد لأن فتحة النسخ المنطقي لـ Postgres تظل المصدر الدائم للحقيقة. إنه wal2json الذي يضمن عدم فقدان أي تغييرات في قاعدة البيانات قبل الاستهلاك، وليس NATS.

JetStream موجود بالفعل في aura-realtime - ولكن لاستخدام مختلف تمامًا. إنه يخدم متجر KV للتواجد عبر المثيلات (aura_presence، مخزن الذاكرة، max_age60 ثانية)، والذي يقوم بمزامنة الأشخاص المتصلين على أي قناة بين جميع النسخ المتماثلة ws-front. حالة مشتركة قصيرة العمر، وليست دفقًا من أحداث مركز السيطرة على الأمراض (CDC) لإعادة تشغيلها.

#
الأمن

إعادة فحص RLS بواسطة المشترك، في كل حدث

لا ينتقل حدث CDC إلى جميع المشتركين في القناة. يقوم المُستقصي بتوجيه كل تغيير إلى أحد العاملين CDC_NUM_SHARDS (4 بشكل افتراضي، مُجزأ بواسطة project_id)، والذي يقوم باستقصاء pg_policiesأولاً. يسمح الجدول الذي لا يحتوي على سياسة RLS لجميع المشتركين - سلوك Supabase - بينما يؤدي الجدول الذي يحتوي على سياسات إلى إجراء فحص لكل مشترك.

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;

تأخذ المعاملة دور Postgres الفعلي للمشترك وتحقن مطالبات JWT الخاصة بها في request.jwt.claims - تلك التي يقرأها auth.uid() وauth.role() على جانب السياسة. ثم يتحقق من أن السطر يظل مرئيًا تحت هذا الدور، ثم يلغي: لا توجد كتابة، قراءة RLS حقيقية. تم التحقق من صحة sub_id فقط في allowed_sub_ids للحدث؛ علامة vec الفارغة تعني عدم وجود أحد، في مغلق بالفشل.

أسوس

يفلت DELETE من هذا الفحص لكل سطر - اختفى الخط، ومن المستحيل إعادة قراءته تحت RLS - ويعود إلى جميع المشتركين في القناة. في المقابل، لا تكشف حمولة الحذف إلا عن أعمدة الهوية (المفتاح الأساسي)، وليس محتوى الصف المحذوف أبدًا.

التفاصيل التي غالبًا ما تعترض طريق الترحيل: wal2json تلتقط فقط المحتوى الكامل لـ UPDATE أو DELETE إذا كان الجدول يحتوي على REPLICA IDENTITY FULL. بدونها، لن يعود سوى INSERT. هذا هو السبب في أن نقطة النهاية التي تتيح الوقت الفعلي على الجدول تقوم بتعيين ALTER TABLE تلقائيًا - وهو إصلاح فعلي للمستودع، وليس مربع اختيار منفصل.

#
تم التحقق من الفخ

لماذا يمكن أن يظل مركز السيطرة على الأمراض (CDC) صامتًا في Kubernetes

تاريخ التسجيل يوثق فخًا حقيقيًا، وليس حالة نظرية. صورة Postgres المستخدمة افتراضيًا في مجموعة Kubernetes (صورة مجتمع المخزون pgvector/pgvector) لا تتضمن المكون الإضافي wal2json. بدونها، يفشل pg_create_logical_replication_slot('aura_cdc_slot', 'wal2json')، ولا يوجد شيء من جانب العميل يشير بشكل صريح إلى ذلك. تظل WebSockets مفتوحة، ويتم قبول الاشتراكات، ولكن لا تحدث أي أحداث لقاعدة البيانات على الإطلاق.

docker/Postgres.Dockerfile (مستخلص حقيقي)dockerfile
RUN apt-get update && \
    apt-get install -y --no-install-recommends \
      postgresql-16-wal2json=2.6-4.pgdg12+1 \
      ...

كان الإصلاح الفعلي هو إنشاء ونشر صورة Postgres مخصصة مع تثبيت هذه الحزمة، ثم توجيه Kubernetes StatefulSet إليها بدلاً من الصورة المخزنة. كما أنه يقوم أيضًا بتعيين wal_level=logicalو max_replication_slots و max_slot_wal_keep_size بشكل صريح كوسيطات لبدء التشغيل - غائبة عن الصور العامة.

من أجل جعل هذا النوع من الفشل قابلاً للملاحظة على وجه التحديد، يعرض cdc-worker علامة cdc_active (خطأ طالما لم يتم تأكيد صحة الفتحة) على نقطة النهاية /health المخصصة لها. هذه إشارة ثنائية يجب مراقبتها - بدلاً من استنتاج مراكز السيطرة على الأمراض (CDC) الميتة من غياب الأحداث من جانب العميل فقط.

#
الحد الحالي

ما لا يغطيه خط الأنابيب هذا بعد: مجموعات مخصصة

تقرأ هذه الآلية متغيرًا واحدًا POSTGRES_REPLICATION_URL، وبالتالي خادم Postgres واحدًا وفتحة منطقية واحدة. يتوافق هذا مع مخطط المستأجرين المتعددين لكل نموذج مشروع (project_<uuid>) في مجموعة مشتركة. في هذا النموذج، كل شيء يعمل: cdc-worker واحد يرى التغييرات في جميع المشاريع في نفس المجموعة، والتوجيه حسب الرسم التخطيطي.

بالنسبة لمشروع على مجموعة CNPG مخصصة - مثيل Postgres المعزول الخاص به - فإن الصورة المستخدمة اليوم (المبنية من صورة مخزون CloudNativePG) تضيف فقط pg_graphql، وليس wal2json. لا شيء أيضًا يقوم بتشغيل cdc-worker لكل مجموعة مخصصة. في الحالة الحالية للتعليمات البرمجية، يظل مركز السيطرة على الأمراض (CDC) الموصوف هنا في الوقت الفعلي يركز على الطبقة المشتركة - وهي حدود معمارية حقيقية، وليست وظيفة خريطة الطريق.

#
جانب العميل

ما لا يتغير: واجهة برمجة تطبيقات postgres_changes الخاصة بـ SDK

تظل هذه الآلية بأكملها غير مرئية من SDK. لا تتغير واجهة برمجة التطبيقات القابلة للتسلسلaurabase-js اعتمادًا على ما إذا كان الحدث يأتي من cdc-worker عبر NATS أو، في التطوير المحلي، من ثنائي 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}'

بالنسبة لبقية سطح SDK المتوافق مع Supabase — المصادقة والتخزين والسياسات RLS —، راجع دليل الترحيل . للتعرف على كيفية تقسيم aura-realtime إلى ثنائيات متعددة في صندوق مساحة عمل واحد، راجع تفاصيل بنية الشحن .

هل أنت جاهز للنشر؟

الواجهة الخلفية الخاصة بك في خمس دقائق.

لا حاجة لبطاقة ائتمان · 500 ميجابايت مجانًا · 50000 MAU