يتيح هذا المحرك دمج ClickHouse مع NATS.
يتيح لك NATS ما يلي:
- نشر subject الرسائل أو الاشتراك فيها.
- معالجة الرسائل الجديدة عند توفرها.
إنشاء جدول
CREATE TABLE [IF NOT EXISTS] [db.]table_name [ON CLUSTER cluster]
(
name1 [type1] [DEFAULT|MATERIALIZED|ALIAS expr1],
name2 [type2] [DEFAULT|MATERIALIZED|ALIAS expr2],
...
) ENGINE = NATS SETTINGS
nats_url = 'host:port',
nats_subjects = 'subject1,subject2,...',
nats_format = 'data_format'[,]
[nats_schema = '',]
[nats_num_consumers = N,]
[nats_queue_group = 'group_name',]
[nats_secure = false,]
[nats_max_reconnect = N,]
[nats_reconnect_wait = N,]
[nats_server_list = 'host1:port1,host2:port2,...',]
[nats_skip_broken_messages = N,]
[nats_max_block_size = N,]
[nats_flush_interval_ms = N,]
[nats_username = 'user',]
[nats_password = 'password',]
[nats_token = 'clickhouse',]
[nats_credentials = '-----BEGIN NATS USER JWT----- ...',]
[nats_startup_connect_tries = 5,]
[nats_max_rows_per_message = 1,]
[nats_commit_on_select = false,]
[nats_handle_error_mode = 'default']المعلمات المطلوبة:
nats_url–host:port(على سبيل المثال،localhost:4222)..nats_subjects– قائمة بالـsubjectالتي يشترك فيها جدول NATS أو ينشر إليها. يدعمsubjectأحرف البدل مثلfoo.*.barأوbaz.>nats_format– تنسيق الرسالة. يستخدم الصياغة نفسها المستخدمة في دالة SQLFORMAT، مثلJSONEachRow. لمزيد من المعلومات، راجع قسم التنسيقات.
المعلمات الاختيارية:
nats_schema– معلمة يجب استخدامها إذا كان التنسيق يتطلب تعريفschema. على سبيل المثال، يتطلب Cap'n Proto مسار ملفschemaواسم الكائن الجذرschema.capnp:Message.nats_stream– اسمstreamموجود في NATS JetStream.nats_consumer_name– اسم durable pull مستهلك موجود في NATS JetStream.nats_num_consumers– عدد الـمستهلكينلكل جدول. القيمة الافتراضية:1. حدِّد عددًا أكبر من الـمستهلكينإذا كان معدل النقل لمستهلك واحد غير كافٍ في NATS core فقط.nats_queue_group– اسم مجموعة queue لمشتركي NATS. القيمة الافتراضية هي اسم الجدول.nats_max_reconnect– مهمل وليس له أي تأثير، إذ تتم إعادة الاتصال بشكل دائم مع مهلةnats_reconnect_wait.nats_reconnect_wait– مقدار الوقت بالمللي ثانية للانتظار بين كل محاولة إعادة اتصال. القيمة الافتراضية:2000.nats_server_list- قائمة الخوادم الخاصة بالاتصال. يمكن تحديدها للاتصال بعنقود NATS.nats_skip_broken_messages- مدى تحمّل محلل رسائل NATS للرسائل غير المتوافقة معschemaلكلblock. القيمة الافتراضية:0. إذا كانnats_skip_broken_messages = Nفسيتخطى المحرك عدد N من رسائل NATS التي يتعذر تحليلها (الرسالة الواحدة تعادل صفًا واحدًا من البيانات).nats_max_block_size- عدد الـrowالتي تجمعها عملية أو عملياتpollلتفريغ البيانات من NATS. القيمة الافتراضية: max_insert_block_size.nats_flush_interval_ms- مهلة تفريغ البيانات المقروءة من NATS. القيمة الافتراضية: stream_flush_interval_ms.nats_wait_for_flush_interval- إذا كانت القيمةtrue، فستبقى دورة البث في الخلفية مفتوحة طوال فترة التفريغ (nats_flush_interval_ms، أوstream_flush_interval_msبخلاف ذلك) بدلًا من الانتهاء بمجرد تفريغ queue المستهلك، مما يتيح تراكم المزيد من الرسائل فيblockواحد مقابل زمن استيعاب إضافي يصل إلى فترة تفريغ واحدة. القيمة الافتراضية:false(سلوك تفريغ فوري منخفض الكمون).nats_username- اسم مستخدم NATS. عند تخزينه في مجموعة مسماة معرّفة في ملف إعدادات الخادم، لا يمكن لـ query تجاوزnats_urlأوnats_server_listالخاصين بالمجموعة.nats_password- كلمة مرور NATS. عند تخزينها في مجموعة مسماة معرّفة في ملف إعدادات الخادم، لا يمكن لـ query تجاوزnats_urlأوnats_server_listالخاصين بالمجموعة.nats_token- رمز مصادقة NATS. عند تخزينه في مجموعة مسماة معرّفة في ملف إعدادات الخادم، لا يمكن لـ query تجاوزnats_urlأوnats_server_listالخاصين بالمجموعة.nats_credential_file- مسار ملف بيانات اعتماد NATS. لا يُقبل إلا من مجموعة مسماة معرّفة في ملف إعدادات الخادم لا يتم تجاوزnats_urlوnats_server_listالخاصين بها بواسطة query، لأن الخادم يفتح المسار بصلاحياته الخاصة. في query، مرّر محتويات الملف فيnats_credentialsبدلًا من ذلك.nats_credentials- محتوى بيانات اعتماد NATS (نفس الحمولة الموجودة في ملف.credsالذي يتضمن JWT المستخدم والبذرة). نظرًا إلى أنه الصيغة الوحيدة التي يمكن لـ query استخدامها، فإنه يستبدلnats_credential_fileالموروث من مجموعة مسماة بدلًا من التعارض معه، ما لم يقفل المشغّل ذلك المسار باستخدام<nats_credential_file overridable="false">. لا يمكن تعيين سلسلة فارغة له لإزالة بيانات الاعتماد التي تحملها مجموعة مسماة.nats_startup_connect_tries- عدد محاولات الاتصال عند بدء التشغيل. القيمة الافتراضية:5.nats_max_rows_per_message— الحد الأقصى لعددrowsالمكتوبة في رسالة NATS واحدة للتنسيقات المعتمدة على الصفوف. (القيمة الافتراضية:1).nats_commit_on_select- ثبّت الرسائل عند تنفيذ query. ينطبق على JetStream فقط؛ إذ لا يحتوي NATS core على إقرارات. القيمة الافتراضية:0.nats_handle_error_mode— كيفية التعامل مع الأخطاء في محرك NATS. القيم الممكنة: default (سيتم طرح الاستثناء إذا فشل تحليل رسالة)، stream (ستُحفَظ رسالة الاستثناء والرسالة الخام فيأعمدة افتراضية_errorو_raw_message).
اتصال SSL:
لإجراء اتصال آمن استخدم nats_secure = 1.
يتم التحكم في التحقق من الشهادة عبر متغير البيئة CLICKHOUSE_NATS_TLS_SECURE.
إذا كانت الشهادة منتهية الصلاحية، أو موقعة ذاتيًا، أو مفقودة، أو غير صالحة لأي سبب آخر، فقم بتعطيل التحقق بتعيين CLICKHOUSE_NATS_TLS_SECURE=0.
الكتابة في جدول NATS:
إذا كان الجدول يقرأ فقط من subject واحد، فسيُنشَر أي إدراج إلى نفس subject.
ومع ذلك، إذا كان الجدول يقرأ من عدة subject، فسنحتاج إلى تحديد الـ subject الذي نرغب في النشر إليه.
لذلك، عند الإدراج في جدول يقرأ من عدة subject، يلزم ضبط stream_like_engine_insert_queue.
يمكنك اختيار أحد الـ subject التي يقرأ منها الجدول ونشر بياناتك هناك. على سبيل المثال:
CREATE TABLE queue (
key UInt64,
value UInt64
) ENGINE = NATS
SETTINGS nats_url = 'localhost:4444',
nats_subjects = 'subject1,subject2',
nats_format = 'JSONEachRow';
INSERT INTO queue
SETTINGS stream_like_engine_insert_queue = 'subject2'
VALUES (1, 1);يمكن أيضًا إضافة إعدادات التنسيق إلى جانب الإعدادات المتعلقة بـ NATS.
مثال:
CREATE TABLE queue (
key UInt64,
value UInt64,
date DateTime
) ENGINE = NATS
SETTINGS nats_url = 'localhost:4444',
nats_subjects = 'subject1',
nats_format = 'JSONEachRow',
date_time_input_format = 'best_effort';يمكن إضافة إعدادات خادم NATS باستخدام ملف إعدادات ClickHouse. وبشكل أكثر تحديدًا، يمكنك إضافة كلمة مرور NATS لمحرك NATS:
<nats>
<user>click</user>
<password>house</password>
<token>clickhouse</token>
</nats>الوصف
لا تكون SELECT مفيدة كثيرًا لقراءة الرسائل (إلا لأغراض استكشاف الأخطاء وإصلاحها)، لأن كل رسالة لا يمكن قراءتها إلا مرة واحدة. والأكثر عملية هو إنشاء مسارات في الوقت الفعلي باستخدام العروض المادية. للقيام بذلك:
- استخدم المحرك لإنشاء مستهلك NATS واعتبره stream بيانات.
- أنشئ جدولًا بالبنية المطلوبة.
- أنشئ عرضًا ماديًا يحوّل البيانات من المحرك ويضعها في جدول أُنشئ مسبقًا.
عندما يرتبط MATERIALIZED VIEW بالمحرك، يبدأ في جمع البيانات في الخلفية. يتيح لك ذلك الاستمرار في تلقي الرسائل من NATS وتحويلها إلى التنسيق المطلوب باستخدام SELECT.
يمكن أن يحتوي جدول NATS واحد على أي عدد تريده من العروض المادية؛ فهي لا تقرأ البيانات من الجدول مباشرةً، بل تستقبل السجلات الجديدة (على شكل كتل)، وبهذه الطريقة يمكنك الكتابة إلى عدة جداول بمستويات مختلفة من التفصيل (مع التجميع وبدونه).
مثال:
CREATE TABLE queue (
key UInt64,
value UInt64
) ENGINE = NATS
SETTINGS nats_url = 'localhost:4444',
nats_subjects = 'subject1',
nats_format = 'JSONEachRow',
date_time_input_format = 'best_effort';
CREATE TABLE daily (key UInt64, value UInt64)
ENGINE = MergeTree() ORDER BY key;
CREATE MATERIALIZED VIEW consumer TO daily
AS SELECT key, value FROM queue;
SELECT key, value FROM daily ORDER BY key;لإيقاف استلام بيانات الـ stream أو لتغيير منطق التحويل، افصل العرض المادي:
DETACH TABLE consumer;
ATTACH TABLE consumer;إذا كنت تريد تغيير الجدول الهدف باستخدام ALTER، فنوصي بتعطيل العرض المادي لتجنّب حدوث اختلافات بين الجدول الهدف والبيانات الواردة من العرض.
الأعمدة الافتراضية
_subject- الـsubject لرسالة NATS. نوع البيانات:String.
أعمدة افتراضية إضافية عندما تكون قيمة nats_handle_error_mode='stream':
_raw_message- الرسالة الخام التي تعذّر تحليلها بنجاح. نوع البيانات:Nullable(String)._error- رسالة الاستثناء التي حدثت أثناء فشل التحليل. نوع البيانات:Nullable(String).
ملاحظة: لا تُملأ الأعمدة الافتراضية _raw_message و _error إلا عند حدوث استثناء أثناء التحليل، وتكون دائمًا NULL عندما يُحلَّلَت الرسالة بنجاح.
دعم تنسيقات البيانات
يدعم محرك NATS جميع التنسيقات التي يدعمها ClickHouse. ويعتمد عدد الصفوف في رسالة NATS الواحدة على ما إذا كان التنسيق يعتمد على الصفوف أم على الكتل:
- بالنسبة إلى التنسيقات المعتمدة على الصفوف، يمكن التحكم في عدد الصفوف في رسالة NATS الواحدة عبر ضبط
nats_max_rows_per_message. - أما بالنسبة إلى التنسيقات المعتمدة على الكتل، فلا يمكن تقسيم الكتلة إلى أجزاء أصغر، لكن يمكن التحكم في عدد الصفوف في الكتلة الواحدة من خلال الإعداد العام max_block_size.
استخدام JetStream
قبل استخدام محرك NATS مع NATS JetStream، يجب إنشاء stream في NATS وdurable pull مستهلك. ويمكنك لهذا الغرض استخدام أداة nats، على سبيل المثال، من حزمة NATS CLI:
إنشاء stream
$ nats stream add
? Stream Name stream_name
? Subjects stream_subject
? Storage file
? Replication 1
? Retention Policy Limits
? Discard Policy Old
? Stream Messages Limit -1
? Per Subject Messages Limit -1
? Total Stream Size -1
? Message TTL -1
? Max Message Size -1
? Duplicate tracking time window 2m0s
? Allow message Roll-ups No
? Allow message deletion Yes
? Allow purging subjects or the entire stream Yes
Stream stream_name was created
Information for Stream stream_name created 2025-10-03 14:12:51
Subjects: stream_subject
Replicas: 1
Storage: File
Options:
Retention: Limits
Acknowledgments: true
Discard Policy: Old
Duplicate Window: 2m0s
Direct Get: true
Allows Msg Delete: true
Allows Purge: true
Allows Per-Message TTL: false
Allows Rollups: false
Limits:
Maximum Messages: unlimited
Maximum Per Subject: unlimited
Maximum Bytes: unlimited
Maximum Age: unlimited
Maximum Message Size: unlimited
Maximum Consumers: unlimited
State:
Messages: 0
Bytes: 0 B
First Sequence: 0
Last Sequence: 0
Active Consumers: 0إنشاء durable pull مستهلك
$ nats consumer add
? Select a Stream stream_name
? Consumer name consumer_name
? Delivery target (empty for Pull Consumers)
? Start policy (all, new, last, subject, 1h, msg sequence) all
? Acknowledgment policy explicit
? Replay policy instant
? Filter Stream by subjects (blank for all)
? Maximum Allowed Deliveries -1
? Maximum Acknowledgments Pending 0
? Deliver headers only without bodies No
? Add a Retry Backoff Policy No
Information for Consumer stream_name > consumer_name created 2025-10-03T14:13:51+03:00
Configuration:
Name: consumer_name
Pull Mode: true
Deliver Policy: All
Ack Policy: Explicit
Ack Wait: 30.00s
Replay Policy: Instant
Max Ack Pending: 1,000
Max Waiting Pulls: 512
State:
Last Delivered Message: Consumer sequence: 0 Stream sequence: 0
Acknowledgment Floor: Consumer sequence: 0 Stream sequence: 0
Outstanding Acks: 0 out of maximum 1,000
Redelivered Messages: 0
Unprocessed Messages: 0
Waiting Pulls: 0 of maximum 512بعد إنشاء stream وdurable pull مستهلك، يمكننا إنشاء جدول باستخدام محرك NATS. وللقيام بذلك، يجب تهيئة: nats_stream وnats_consumer_name وnats_subjects:
CREATE TABLE nats_jet_stream (
key UInt64,
value UInt64
) ENGINE NATS
SETTINGS nats_url = 'localhost:4222',
nats_stream = 'stream_name',
nats_consumer_name = 'consumer_name',
nats_subjects = 'stream_subject',
nats_format = 'JSONEachRow';توفر جداول JetStream تسليمًا مرة واحدة على الأقل: لا يتم الإقرار بالرسالة إلا بعد إدراجها في العروض المادية التابعة، لذا تظل الرسالة التي يفشل إدراجها أو ينقطع إدراجها من دون إقرار، وتُعاد تسليمها. لا يوفر Core NATS (من دون JetStream) إقرارًا أو إعادة تشغيل، لذا فهو يضمن التسليم مرة واحدة على الأكثر، وتُفقد الرسالة التي ينقطع تسليمها.