الانتقال إلى المحتوى الرئيسي
هذا هو موصل sink الرسمي لـ Apache Flink والمدعوم من ClickHouse. وقد بُني باستخدام AsyncSinkBase في Flink وJava client الرسمي لـ ClickHouse. يدعم هذا الموصل واجهة DataStream API في Apache Flink. ومن المخطط دعم Table API في إصدار لاحق.

المتطلبات

  • Java 11+ (لـ Flink 1.17+) أو 17+ (لـ Flink 2.0+)
  • Apache Flink 1.17+
ينقسم الموصل إلى حزمتَي artifact لدعم كلٍّ من Flink 1.17+ وFlink 2.0+. اختر الـ artifact المطابق لإصدار Flink الذي تريده:
لم يُختبر هذا الموصل مع إصدارات Flink الأقدم من 1.17.2

التثبيت والإعداد

استيراده كتبعية

تنزيل ملف JAR الثنائي

نمط تسمية ملف JAR الثنائي هو:
حيث: يمكنك العثور على جميع ملفات JAR المُتاحة والصادرة في Maven Central Repository.

استخدام واجهة برمجة تطبيقات DataStream

مقتطف

لنفترض أنك تريد إدراج بيانات CSV خام في ClickHouse:
يمكنك العثور على المزيد من الأمثلة والمقتطفات في اختباراتنا:

مثال للبدء السريع

أنشأنا مثالًا قائمًا على Maven لتسهيل البدء باستخدام ClickHouse Sink: للحصول على إرشادات أكثر تفصيلًا، راجع دليل الأمثلة

خيارات الاتصال بواجهة برمجة تطبيقات DataStream

خيارات عميل ClickHouse

يجب تمرير options وserverSettings إلى العميل بصيغة Map<String, String>. وسيؤدي استخدام خريطة فارغة لأيٍّ منهما إلى استخدام الإعدادات الافتراضية للعميل أو الخادم، على الترتيب.
جميع خيارات عميل Java المتاحة مُدرجة في ClientConfigProperties.java وصفحة التوثيق هذه.جميع إعدادات جلسة الخادم المتاحة مُدرجة في صفحة التوثيق هذه.
على سبيل المثال:

خيارات الـ sink

تأتي الخيارات التالية مباشرةً من AsyncSinkBase في Flink:

أنواع البيانات المدعومة

يوفّر الجدول أدناه مرجعًا سريعًا لتحويل أنواع البيانات عند الإدراج من Flink إلى ClickHouse. ملاحظات:
  • يجب توفير ZoneId عند إجراء عمليات على التاريخ.
  • يجب توفير الدقة والمقياس عند إجراء عمليات على القيم العشرية.
  • لكي يتمكن ClickHouse من تحليل String في Java على أنه JSON، يجب تمكين enableJsonSupportAsString في ClickHouseClientConfig.
  • يتطلب الموصّل ElementConvertor لربط العناصر في DataStream المدخل بحمولات ClickHouse. ولهذا الغرض، يوفّر الموصّل ClickHouseConvertor وPOJOConvertor، ويمكنك استخدامهما لتنفيذ هذا الربط باستخدام طرق التسلسل الخاصة بـ DataWriter المذكورة أعلاه.

تنسيقات الإدخال المدعومة

يمكنك العثور على قائمة تنسيقات الإدخال المتاحة في ClickHouse في صفحة التوثيق هذه وClickHouseFormat.java. لتحديد التنسيق الذي ينبغي أن يستخدمه الموصل لتحويل DataStream إلى ClickHouse حمولات، استخدم الدالة setClickHouseFormat. على سبيل المثال:
افتراضيًا، سيستخدم الموصل تنسيق RowBinaryWithDefaults أو RowBinary إذا ضُبطت القيمة setSupportDefault في ClickHouseClientConfig صراحةً على true أو false، على الترتيب.

المقاييس

يوفّر الموصل المقاييس الإضافية التالية بالإضافة إلى مقاييس Flink الحالية:

القيود

  • يوفّر sink حاليًا ضمان تسليم مرة واحدة على الأقل. ويجري تتبّع العمل لتحقيق exactly-once semantics هنا.
  • لا يدعم sink بعدُ قائمة انتظار الرسائل الميتة (DLQ) لتخزين السجلات غير القابلة للمعالجة مؤقتًا. وحتى ذلك الحين، سيحاول الموصل إعادة إدراج السجلات التي تفشل، وسيتخلّص منها إذا لم ينجح ذلك. ويجري تتبّع هذه الميزة هنا.
  • لا يدعم sink بعدُ الإنشاء عبر Table API الخاصة بـ Flink أو Flink SQL. ويجري تتبّع هذه الميزة هنا.

توافق إصدارات ClickHouse والأمان

  • يُختبَر موصل مع مجموعة من إصدارات ClickHouse الحديثة، بما في ذلك latest وhead، عبر سير عمل CI يومي. وتُحدَّث الإصدارات المختبَرة دوريًا مع اعتماد إصدارات ClickHouse الجديدة. اطّلع هنا على الإصدارات التي يُختبَر موصل عليها يوميًا.
  • راجع سياسة أمان ClickHouse للتعرّف على الثغرات الأمنية المعروفة وكيفية الإبلاغ عن أي ثغرة.
  • نوصي بترقية موصل باستمرار حتى لا تفوتك الإصلاحات الأمنية والتحسينات الجديدة.
  • إذا واجهت مشكلة في الترحيل، فيُرجى إنشاء issue على GitHub وسنرد عليك!
  • للحصول على أفضل أداء، تأكد من أن نوع العنصر في DataStream لديك ليس من النوع Generic - راجع هذا الشرح لتمييز الأنواع في Flink. فالعناصر غير العامة تتجنب كلفة التسلسل الإضافية التي يفرضها Kryo وتُحسّن معدل النقل إلى ClickHouse.
  • نوصي بضبط maxBatchSize على 1000 كحد أدنى، ويفضَّل أن يكون بين 10,000 و100,000. راجع هذا الدليل حول عمليات الإدراج المجمّعة لمزيد من المعلومات.
  • لإجراء إزالة التكرار أو upsert إلى ClickHouse بأسلوب OLTP، راجع صفحة التوثيق هذه. ملاحظة: لا تخلط بين هذا وبين إزالة التكرار على مستوى الدُفعات التي تحدث عند إعادة المحاولة.

استكشاف الأخطاء وإصلاحها

CANNOT_READ_ALL_DATA

قد يظهر الخطأ التالي:
السبب: في معظم الحالات، يعني الخطأ CANNOT_READ_ALL_DATA أن مخطط جدول ClickHouse لديك لم يعد متطابقًا مع مخطط سجل Flink. وقد يحدث ذلك عندما يُعدَّل أحدهما أو كلاهما بطريقة غير متوافقة مع الإصدارات السابقة. الحل: حدِّث المخطط في جدول ClickHouse لديك أو نوع بيانات إدخال الموصل (أو كليهما) بحيث يصبحان متوافقين. وإذا لزم الأمر، فارجع إلى تعيين الأنواع لمعرفة كيفية ربط أنواع Java بأنواع ClickHouse. ملاحظة: إذا كانت لا تزال هناك سجلات قيد النقل، فستحتاج إلى إعادة تعيين حالة Flink عند إعادة تشغيل الموصل.

انخفاض معدل النقل

قد تلاحظ أن معدل نقل الموصل لا يزداد بما يتناسب مع توازي المهمة (عدد مهام Flink) عند الكتابة إلى ClickHouse. السبب: قد تؤدي عملية دمج الأجزاء في الخلفية في ClickHouse إلى إبطاء عمليات الإدراج. يمكن أن يحدث ذلك عندما يكون حجم الدفعة المُعدّ صغيرًا جدًا، أو عندما يقوم الموصل بعملية التفريغ بشكل متكرر جدًا، أو بسبب الجمع بين الأمرين. الحل: راقب المقياسين numRequestSubmitted وactualRecordsPerBatch للمساعدة في تحديد كيفية ضبط حجم الدفعة (maxBatchSize) وعدد مرات التفريغ. راجع أيضًا الاستخدام المتقدم والموصى به للاطلاع على توصيات بشأن حجم الدفعات.

هناك صفوف مفقودة في جدول ClickHouse الخاص بي

السبب: أُسقِطت الدفعة (أو الدفعات) إما بسبب فشل غير قابل لإعادة المحاولة، أو لعدم إمكانية إدراجها ضمن عدد محاولات إعادة المحاولة المُعدّ (يمكن ضبطه عبر ClickHouseClientConfig.setNumberOfRetries()). ملاحظة: يحاول الموصل، افتراضيًا، إعادة إدراج الدفعة حتى 3 مرات قبل إسقاطها. الحل: افحص سجلات TaskManager و/أو تتبّع المكدس لتحديد السبب الجذري.

المساهمة والدعم

إذا كنت ترغب في المساهمة في المشروع أو الإبلاغ عن أي مشكلات، فنحن نرحّب بمشاركتك! تفضّل بزيارة مستودع GitHub الخاص بنا لفتح issue أو اقتراح تحسينات أو إرسال pull request. نرحّب بالمساهمات! يُرجى الاطلاع على دليل المساهمة في المستودع قبل البدء. شكرًا لمساعدتك في تحسين موصل ClickHouse لـ Flink!
آخر تعديل في ٢ يوليو ٢٠٢٦