من Kafka إلى ClickHouse
إذا كنت تستخدم ClickHouse Cloud، فننصحك باستخدام ClickPipes بدلاً منه. يوفّر ClickPipes دعمًا أصليًا لاتصالات الشبكة الخاصة، ولتوسيع موارد الإدخال وموارد العنقود بشكل مستقل، ولمراقبة شاملة لإدخال بيانات Kafka المتدفقة إلى ClickHouse.
نظرة عامة
TO وجهة البيانات، والتي تكون عادةً جدولًا من عائلة MergeTree. يوضّح الشكل التالي هذه العملية:
الخطوات
1
التحضير
إذا كانت لديك بيانات مُعبّأة في topic الهدف، فيمكنك تكييف ما يلي لاستخدامه في مجموعة بياناتك. وبدلًا من ذلك، تتوفر مجموعة بيانات GitHub تجريبية هنا. تُستخدم مجموعة البيانات هذه في الأمثلة أدناه، وتعتمد مخطط GitHub المبسّط وجزءًا فقط من الصفوف (وتحديدًا، نقتصر على أحداث GitHub المتعلقة بـمستودع ClickHouse)، مقارنةً بمجموعة البيانات الكاملة المتاحة هنا، اختصارًا. ومع ذلك، يظل هذا كافيًا لكي تعمل معظم الاستعلامات المنشورة مع مجموعة البيانات.
2
هيّئ ClickHouse
هذه الخطوة مطلوبة إذا كنت تتصل بـ Kafka آمن. لا يمكن تمرير هذه الإعدادات عبر أوامر SQL DDL، ويجب تهيئتها في ملف ClickHouse config.xml. نفترض أنك تتصل بمثيل مؤمّن باستخدام SASL. وهذه هي أبسط طريقة عند التفاعل مع Confluent Cloud.إما أن تضع المقتطف أعلاه في ملف جديد داخل دليل بمجرد إنشاء قاعدة البيانات، ستحتاج إلى الانتقال إليها:
conf.d/ لديك أو تدمجه في ملفات التهيئة الحالية. للاطلاع على الإعدادات التي يمكن تهيئتها، انظر هنا.سننشئ أيضًا قاعدة بيانات باسم KafkaEngine لاستخدامها في هذا الدليل العملي:3
أنشئ جدول الوجهة
جهّز جدول الوجهة. في المثال أدناه، نستخدم مخطط GitHub المبسّط اختصارًا. لاحظ أنه رغم أننا نستخدم محرك الجدول MergeTree، يمكن تكييف هذا المثال بسهولة للعمل مع أي محرك ضمن عائلة MergeTree.
4
أنشئ الـ topic واملأه
بعد ذلك، سننشئ topic. هناك عدة أدوات يمكننا استخدامها لهذا الغرض. إذا كنا نشغّل Kafka محليًا على أجهزتنا أو داخل حاوية Docker، فإن RPK يعمل بكفاءة. يمكننا إنشاء topic باسم إذا كنّا نشغّل Kafka على Confluent Cloud، فقد نفضّل استخدام Confluent CLI:نحتاج الآن إلى إرسال بعض البيانات إلى هذا الـ topic، وسنفعل ذلك باستخدام kcat. يمكننا تشغيل أمر مشابه لما يلي إذا كنا نشغّل Kafka محليًا مع تعطيل المصادقة:أو ما يلي إذا كان عنقود Kafka لدينا يستخدم SASL للمصادقة:تحتوي مجموعة البيانات على 200,000 صف، لذا يُفترض إدخالها خلال بضع ثوانٍ فقط. إذا كنت تريد العمل مع مجموعة بيانات أكبر، فألقِ نظرة على قسم مجموعات البيانات الكبيرة في مستودع GitHub ClickHouse/kafka-samples.
github مع 5 partitions عبر تشغيل الأمر التالي:5
أنشئ محرك الجدول Kafka
ينشئ المثال أدناه محرك جدول له المخطط نفسه لجدول MergeTree. وليس هذا ضروريًا تمامًا، إذ يمكن أن يكون لديك اسم مستعار أو أعمدة مؤقتة في جدول الوجهة. لكن الإعدادات مهمة — لاحظ استخدام نناقش أدناه إعدادات المحرك وضبط الأداء. عند هذه المرحلة، ينبغي أن يقرأ استعلام select بسيط على الجدول
JSONEachRow بوصفه نوع البيانات لاستهلاك JSON من موضوع Kafka. وتمثّل القيمتان github وclickhouse اسم الموضوع واسم مجموعة المستهلكين، على الترتيب. ويمكن أن تكون الموضوعات في الواقع قائمة من القيم.github_queue بعض الصفوف. لاحظ أن هذا سيؤدي إلى تقديم offsets الخاصة بالمستهلك، مما يمنع إعادة قراءة هذه الصفوف دون إعادة تعيين. لاحظ أيضًا قيمة limit والمعلمة المطلوبة stream_like_engine_allow_direct_select.6
أنشئ العرض المادي
سيربط العرض المادي بين الجدولين اللذين أُنشئا سابقًا، إذ يقرأ البيانات من محرك الجدول Kafka ويُدرجها في جدول MergeTree الهدف. يمكننا إجراء عدد من تحويلات البيانات. سنجري عملية قراءة وإدراج بسيطة. ويفترض استخدام * أن أسماء الأعمدة متطابقة تمامًا (مع مراعاة حالة الأحرف).عند إنشائه، يتصل العرض المادي بمحرك Kafka ويبدأ القراءة، مُدرجًا الصفوف في جدول الوجهة. وستستمر هذه العملية بلا توقف، مع استهلاك أي رسائل جديدة تُدرج في Kafka. يمكنك إعادة تشغيل برنامج الإدراج النصي لإدراج المزيد من الرسائل في Kafka.
7
تأكد من إدراج الصفوف
تأكد من أن البيانات موجودة في جدول الوجهة:يُفترض أن ترى 200,000 صفًا:
العمليات الشائعة
إيقاف استهلاك الرسائل & إعادة تشغيله
إضافة البيانات الوصفية لـ Kafka
_.
يمكن العثور على قائمة كاملة بالأعمدة الافتراضية هنا.
لتحديث جدولنا بالأعمدة الافتراضية، سنحتاج إلى حذف العرض المادي، ثم إعادة إرفاق جدول محرك Kafka، ثم إعادة إنشاء العرض المادي.
تعديل إعدادات محرك Kafka
استكشاف المشكلات وإصلاحها
التعامل مع الرسائل غير السليمة
- تعامل مع حقل الرسالة على أنه سلاسل نصية. ويمكن استخدام الدوال في عبارة
materialized viewلإجراء التنظيف وCAST عند الحاجة. لا ينبغي أن يُعدّ هذا حلًا لبيئة production، لكنه قد يساعد في الإدخال لمرة واحدة. - إذا كنت تستهلك JSON من topic، باستخدام التنسيق JSONEachRow، فاستخدم الإعداد
input_format_skip_unknown_fields. عند كتابة البيانات، يرفع ClickHouse افتراضيًا استثناءً إذا كانت بيانات الإدخال تحتوي على أعمدة غير موجودة في الجدول الهدف. ومع ذلك، إذا كان هذا الخيار مفعّلًا، فسيتم تجاهل هذه الأعمدة الزائدة. ومرة أخرى، لا يُعدّ هذا حلًا على مستوى production وقد يسبب التباسًا للآخرين. - ضع في اعتبارك الإعداد
kafka_skip_broken_messages. يتطلب هذا من المستخدم تحديد مستوى التحمّل لكل كتلة تجاه الرسائل غير السليمة، وذلك في سياق kafka_max_block_size. وإذا تم تجاوز هذا التحمّل (ويُقاس بعدد الرسائل المطلق)، فسيعود سلوك الاستثناء المعتاد، وسيتم تخطي الرسائل الأخرى.
دلالات التسليم والتحديات المرتبطة بالتكرارات
عمليات الإدراج القائمة على النصاب
ClickHouse إلى Kafka
الخطوات
1
إدراج الصفوف مباشرةً
أولًا، تحقّق من عدد صفوف الجدول الهدف.يُفترض أن يكون لديك 200,000 صف:أدرِج الآن صفوفًا من الجدول الهدف GitHub مرة أخرى في محرك جدول Kafka github_queue. لاحظ كيف نستخدم تنسيق JSONEachRow وكيف نستخدم LIMIT لتقييد استعلام SELECT إلى 100.أعِد عدّ الصفوف في GitHub للتأكد من أنها قد زادت بمقدار 100. وكما هو موضح في المخطط أعلاه، أُدرِجت الصفوف في Kafka عبر محرك جدول Kafka قبل أن يعيد المحرّك نفسه قراءتها ويُدرجها عرضنا المادي في جدول GitHub الهدف!يُفترض أن ترى 100 صف إضافيًا:
2
استخدام العروض المادية
يمكننا استخدام العروض المادية لإرسال الرسائل إلى محرك Kafka (وإلى topic) عند إدراج المستندات في جدول. وعند إدراج صفوف في جدول GitHub، يتم تشغيل عرض مادي، ما يؤدي إلى إدراج الصفوف مرة أخرى في محرك Kafka وفي topic جديد. ويوضح الشكل التالي ذلك بشكل أفضل:أنشئ Kafka topic جديدًا باسم أنشئ الآن عرضًا ماديًا جديدًا إذا أجريت عملية insert في topic الأصلي المسمى github، والذي أُنشئ كجزء من Kafka to ClickHouse، فستظهر المستندات تلقائيًا في topic “github_clickhouse”. أكِّد ذلك باستخدام أدوات Kafka الأصلية. على سبيل المثال، نُدرج أدناه 100 صف في topic github باستخدام kcat لـ topic مستضاف على Confluent Cloud:من المفترض أن تؤكد قراءة من topic على الرغم من كون هذا المثال تفصيليًا، إلا أنه يوضح قوة العروض المادية عند استخدامها بالتزامن مع محرك Kafka.
github_out أو ما يعادله. وتأكد من أن محرك جدول Kafka باسم github_out_queue يشير إلى هذا الـ topic.github_out_mv يشير إلى جدول GitHub، بحيث يُدرِج الصفوف في المحرّك المذكور أعلاه عند تفعيله. ونتيجةً لذلك، ستُدفَع الإضافات إلى جدول GitHub إلى Kafka topic الجديد لدينا.github_out تسليم الرسائل.العناقيد والأداء
العمل مع عنقود ClickHouse
ضبط الأداء
- يختلف الأداء بحسب حجم الرسائل، والتنسيق، وأنواع الجداول المستهدفة. ويمكن اعتبار 100 ألف صف/ثانية على table engine واحد معدلًا يمكن تحقيقه. افتراضيًا، تُقرأ الرسائل على شكل كتل، ويتحكم في ذلك المعلَمة kafka_max_block_size. افتراضيًا، تُضبط هذه القيمة على max_insert_block_size، وقيمتها الافتراضية 1,048,576. وما لم تكن الرسائل كبيرة جدًا، فمن الأفضل تقريبًا دائمًا زيادة هذه القيمة. كما أن القيم بين 500 ألف ومليون ليست غير شائعة. اختبر هذا وقيّم أثره في أداء معدل النقل.
- يمكن زيادة عدد المستهلكين لمحرك الجدول باستخدام kafka_num_consumers. لكن افتراضيًا، ستُسلسل عمليات الإدراج في thread واحد ما لم يتم تغيير kafka_thread_per_consumer عن القيمة الافتراضية 1. اضبطه على 1 لضمان تنفيذ عمليات التفريغ بالتوازي. لاحظ أن إنشاء جدول Kafka engine بعدد N من المستهلكين (ومع kafka_thread_per_consumer=1) يكافئ منطقيًا إنشاء N من Kafka engines، لكل منها materialized view وkafka_thread_per_consumer=0.
- زيادة عدد المستهلكين ليست بلا تكلفة. إذ يحتفظ كل مستهلك بمخازنه المؤقتة وthreads الخاصة به، مما يزيد الحمل الإضافي على الخادم. لذلك، انتبه إلى الحمل الإضافي الناتج عن المستهلكين، ووسّع أفقيًا عبر عنقود أولًا إن أمكن.
- إذا كان معدل نقل رسائل Kafka متغيرًا وكانت التأخيرات مقبولة، ففكّر في زيادة stream_flush_interval_ms لضمان تفريغ كتل أكبر.
- يحدّد background_message_broker_schedule_pool_size عدد threads التي تنفّذ tasks في الخلفية. وتُستخدم هذه threads في Kafka streaming. يُطبَّق هذا الإعداد عند بدء تشغيل ClickHouse server ولا يمكن تغييره ضمن session للمستخدم، وقيمته الافتراضية 16. إذا رأيت حالات timeout في السجلات، فقد يكون من المناسب زيادة هذه القيمة.
- في الاتصال مع Kafka، تُستخدم مكتبة librdkafka، وهي بدورها تنشئ threads. لذلك، قد يؤدي العدد الكبير من جداول Kafka أو المستهلكين إلى عدد كبير من عمليات تبديل السياق. وزّع هذا الحمل عبر عنقود، مع الاكتفاء بتكرار الجداول المستهدفة فقط إن أمكن، أو فكّر في استخدام محرك جدول للقراءة من topics متعددة، إذ إن قائمة من القيم مدعومة. كما يمكن قراءة عدة materialized views من جدول واحد، بحيث يرشّح كلٌّ منها البيانات الخاصة بـ topic معيّن.
إعدادات إضافية
- Kafka_max_wait_ms - مدة الانتظار بالمللي ثانية لقراءة الرسائل من Kafka قبل إعادة المحاولة. يُضبط هذا الإعداد على مستوى profile المستخدم، وتكون قيمته الافتراضية 5000.