Skip to main content
قبل ظهور مزود ClickHouse الرسمي، كان معظم مستخدمي Airflow يتصلون بـ ClickHouse عبر حزمة المجتمع airflow-clickhouse-plugin. يستعرض هذا الدليل خطوات ترحيل عملية نشر قائمة إلى apache-airflow-providers-clickhousedb. ولمعرفة كيفية عمل المزود نفسه، راجع ربط Apache Airflow بـ ClickHouse.

لماذا مزوّد أصلي؟

المزوّدات (Providers) هي الطريقة التي يتكامل بها Airflow مع الأنظمة الخارجية. ويتم إصدارها واختبارها وتوثيقها مع بقية منظومة Airflow، وهذا المزوّد يتولى صيانته مجتمع Airflow بالتعاون مع ClickHouse team. وهو مبني على clickhouse-connect، وهو عميل بايثون الذي تطوّره ClickHouse وتدعمه بنفسها، لا على driver يصونه المجتمع. وبذلك تصل ميزات الخادم الجديدة وإصلاحاته إلى مستخدمي Airflow عبر مسار مدعوم. ويمنحك الانتقال إلى المزوّد package له موطن رسمي، وعوامل common.sql القياسية ومستشعراتها، ونوع اتصال يظهر في UI الخاص بـ Airflow كما هو الحال مع أي database أخرى. ولا يقتصر الاختلاف بين الحزمتين على مسارات import. فالـ plugin يتواصل مع ClickHouse عبر بروتوكول TCP الأصلي باستخدام clickhouse-driver، أما المزوّد فيتواصل عبر HTTP(S) باستخدام clickhouse-connect ويندمج مع عوامل common.sql العامة بدلاً من توفير عوامل مخصصة لـ ClickHouse.
اقرأ الدليل بالكامل مرة واحدة قبل تغيير أي شيء. فتغيير الاتصال على وجه الخصوص يؤثر على كل DAG في الوقت نفسه.

لمحة سريعة

الخطوة 1: التحقق من المتطلبات المسبقة والتثبيت

يتطلب المزوّد إصدار Airflow 2.11 أو أحدث، والإصدار 1.32.0 أو أحدث من apache-airflow-providers-common-sql. قم بترقية Airflow أولاً إذا كنت تستخدم إصداراً أقدم.
توجد الحزمتان في مساحتي أسماء مختلفتين في بايثون، لذا يمكن تثبيتهما جنبًا إلى جنب أثناء ترحيل كل DAG على حدة. أزل الـ plugin بعد أن لا يعود أي شيء يستورده:

الخطوة 2: تحديث الاتصالات

هذه هي الخطوة التي يؤدي تجاوزها إلى تعطّل الأمور. فكل اتصال ClickHouse قائم يشير إلى المنفذ native، بينما يحتاج provider إلى منفذ HTTP. إذا كان ClickHouse خلف firewall أو load balancer، فتأكّد من إمكانية الوصول إلى منفذ HTTP من الـ workers قبل التبديل. وتحقّق من أن HTTP interface مُمكّن على الخادم (http_port أو https_port في تهيئة الخادم). أما ClickHouse Cloud فلا يعرض HTTPS إلا على المنفذ 8443. أما الاتصالات المخزّنة كعناوين URI (clickhouse://user:pass@host:9000/db?secure=true) فهي تحمل بالفعل نوع الاتصال clickhouse لأن Airflow يستنتجه من الـ scheme في عنوان URI. ولا يتغيّر فيها سوى المنفذ والـ keys الإضافية؛ وتُحلَّل قيم سلسلة الاستعلام على أنها JSON، لذا تبقى secure=true قيمة منطقية.

إضافات الاتصال

كان الـ plugin يمرّر كل مفتاح في extra مباشرةً إلى clickhouse_driver.Client، لذا قد تحمل الاتصالات أي وسيطات كلمات مفتاحية خاصة بـ clickhouse-driver. أما الـ provider فلا يقرأ سوى مجموعة ثابتة من المفاتيح ويمرّر ما عداها عبر client_kwargs. حوّلها على النحو التالي: قبل:
بعد:
تحقّق من كل اتصال مُرحَّل قبل تعديل DAGs:

الخطوة 3: استبدال عمليات الاستيراد

لم يكن الغرض من الـ wrappers الخاصة بـ common.sql والمسبوقة بـ ClickHouse سوى حقن الخطاف الخاص بالـ plugin. أما الـ provider فيسجّل نوع الاتصال clickhouse، ومن ثم تستطيع أصناف common.sql غير المسبوقة تحديد الخطاف من الاتصال بنفسها. فإذا كنت تستخدم تلك الـ wrappers، يقتصر الترحيل عادةً على تعديل سطر الـ import، وحذف السابقة ClickHouse، وتمرير conn_id بشكل صريح (انظر الخطوة 7).

الخطوة 4: من ClickHouseOperator إلى SQLExecuteQueryOperator

قبل:
بعد:

نتائج العبارات المتعددة

كان الـ plugin يدفع نتيجة آخر عبارة إلى XCom. أما SQLExecuteQueryOperator فيُعيد نتيجة واحدة لكل عبارة عندما يكون sql قائمة، لذا فإن المثال أعلاه يدفع [[], [(12345.0,)]] بينما كان الـ plugin يدفع [(12345.0,)]. اختر أحد الخيارين التاليين:
  • غيّر xcom_pull في المرحلة اللاحقة ليأخذ العنصر الأخير.
  • مرّر العبارات كسلسلة نصية واحدة مفصولة بـ ; واضبط split_statements=True. ومع القيمة الافتراضية return_last=True يدفع المُشغّل عندئذٍ صفوف آخر عبارة فقط، بما يطابق سلوك الـ plugin.
إن ClickHouseOperator الذي كان يُدرج قائمة من الصفوف عبر parameters يصبح SQLInsertRowsOperator:
مرّر دائماً columns؛ فبدونها يبحث الـ operator عن الجدول عبر SQLAlchemy.

الحفاظ على column types

كانت with_column_types=True تُعيد (rows, [(name, type), ...]). يمكن تحقيق الأمر نفسه باستخدام handler، حيث يعرض مؤشر clickhouse-connect أسماء أنواع ClickHouse في cursor.description:

الخطوة 5: من ClickHouseHook.execute إلى طرق DbApiHook

كان خطاف الـ plugin يوفّر طريقة واحدة، execute، تُحاكي clickhouse_driver.Client.execute. أما خطاف الـ provider فهو DbApiHook، ولذلك يحصل على الطرق القياسية المتوفرة في كل provider SQL آخر: run وget_records وget_first وget_pandas_df وget_df وinsert_rows وtest_connection. ولم تتغيّر وسيطتا المُنشئ clickhouse_conn_id وdatabase. يُستورد fetch_all_handler وبقية الـ handlers من airflow.providers.common.sql.hooks.handlers. أكثر أنماط استخدام الخطاف شيوعًا في شيفرة عصر الـ plugin هو الإدراج بالجملة (bulk insert). ولا يمكن تنفيذه باستدعاء run بسيط، لأن مؤشر DB-API سيحاول تضمين الصفوف داخل سلسلة SQL. استخدم بدلًا من ذلك الإدراج الأصلي (native insert): قبل:
بعد:
تتطلب bulk_insert_rows تمرير column_names. أما batch_size فهي اختيارية وتحدّ من استهلاك الذاكرة مع المدخلات الضخمة جدًا. كذلك تنتهي الدالة العامة insert_rows(table, rows, target_fields=[...], executemany=True) بعملية إدخال أصلية (native insert)، لكن فقط عند استخدام executemany=True؛ إذ يرسل السلوك الافتراضي طلب HTTP واحدًا لكل صف. ولكل ما لا تغطيه واجهة DB-API، تُعيد get_client() عميل clickhouse-connect الخام المُهيّأ انطلاقًا من اتصال Airflow، وهو البديل لكل argument خاص بـ clickhouse-driver كان الـ plugin يوفّره:

الخطوة 6: من ClickHouseSensor إلى SqlSensor

هذا هو الاستبدال الوحيد الذي يتغيّر فيه دخل الدالة القابلة للاستدعاء. ولأن الـ plugin كان يمرّر الـ result set كاملًا، فإن شيفرة المستشعر القائمة تعتمد على الفهرسة داخله. أزِل هذه الفهرسة عند الترحيل: قبل:
بعد:
إذا كان الكائن القابل للاستدعاء يحتاج إلى الصف بأكمله، فمرّر selector=lambda row: row. وإذا كان يحتاج إلى جميع الصفوف، فاكتب الـ check في SQL بحيث يُعيد الاستعلام قيمة منطقية واحدة أو عدداً.

الخطوة 7: عائلة الأغلفة common.sql

الشيفرة التي استخدمت ClickHouseSQLExecuteQueryOperator وClickHouseSqlSensor وبقية الأغلفة المسبوقة بـ ClickHouse تحتاج إلى أقل قدر من العمل:
  • غيّر الاستيراد إلى وحدة common.sql واحذف السابقة ClickHouse من اسم الصنف.
  • مرّر conn_id بشكل صريح. كانت الأغلفة تتعامل مع conn_id الغائب أو None على أنه clickhouse_default، أما أصناف common.sql فليست لها قيمة افتراضية وتفشل بدونه. ويكفي default_args={"conn_id": "clickhouse_default"} لتغطية DAG بأكمله.
  • لا يزال database= في المشغّلات وhook_params={"schema": ...} في المستشعر يعملان، إذ يتعامل خطاف المزوّد مع schema كاسم بديل لـ database.
  • يصبح ClickHouseDbApiHook هو ClickHouseHook. ولا يزال وسيط الباني schema مقبولًا كاسم بديل، أما database فهي الصيغة الأصلية في ClickHouse ولها الأسبقية عند تمرير الاثنين معًا.
  • لا يزال الاتصال بحاجة إلى تغييرات المنفذ والإضافات الواردة في الخطوة 2، فالأغلفة كانت تستخدم البروتوكول الأصلي أيضًا.

اختلافات السلوك التي يجب مراجعتها

حتى بعد نجاح ترجمة الشيفرة، تختلف بعض السلوكيات في وقت التشغيل. قيمة XCom لمهام INSERT. كان الـ plugin يدفع ما يُرجعه clickhouse-driver، وهو في حالة الإدخال بصيغة VALUES مع معاملات عدد الصفوف المُدخلة. أما الـ provider فيدفع مجموعة نتائج فارغة للعبارات التي لا تُرجع صفوفًا. لذا على المهام اللاحقة التي تقرأ عدد الصفوف من XCom أن تحصل عليه بطريقة أخرى، مثل تنفيذ SELECT count() لاحق. إعدادات الجلسة مقابل عبارات SET. تُنفّذ كلتا الحزمتين قائمة عبارات متعددة عبر اتصال واحد، وclickhouse-connect ينشئ جلسة لكل عميل افتراضيًا، لذا فإن عبارة SET في بداية القائمة ينبغي أن تنطبق على العبارات اللاحقة. ومع ذلك، يُفضّل استخدام session_settings؛ فهو صريح وقابل للقولبة، ويعمل بالطريقة ذاتها سواء حافظ الخادم أو أي proxy وسيط على الجلسة أم لا. تحقّق من السلوك في بيئتك إذا كانت مخططات DAG لديك تعتمد على SET. الاستثناءات. أصبحت الأخطاء الآن clickhouse_connect.driver.exceptions.DatabaseError أو OperationalError أو ProgrammingError بدلًا من clickhouse_driver.errors.ServerException وNetworkError. حدّث عبارات except وشيفرة on_failure_callback ومنطق إعادة المحاولة الذي يفحص أنواع الاستثناءات. الضغط. كان clickhouse-driver يُبقي الضغط معطّلًا ما لم يُضبط compression؛ أما clickhouse-connect فيُمكّن ضغط استجابات HTTP افتراضيًا ويتفاوض على الخوارزمية مع الخادم. اضبط "compress": false في حقل extra الخاص بالاتصال لاستعادة السلوك القديم. ولم تعد حزمة clickhouse-cityhash التي كان الـ plugin يحتاجها للضغط الأصلي مطلوبة؛ بينما تبقى lz4 مثبّتة كتابعة لـ clickhouse-connect. تعيين الأنواع. يُرجع كلا الـ driver أنواع Python أصلية، لكنهما قاعدتا شيفرة مختلفتان. راجع المهام التي تعتمد على أنواع محددة بدقة مع DateTime64 ذات المناطق الزمنية، وDecimal، وUUID، وcolumn من نوع Nullable، وقيم Array أو Map المتداخلة، خصوصًا حين تُدفع النتيجة إلى XCom وتُستهلك في مراحل لاحقة. تعريف الاستعلامات في system.query_log. تصل الاستعلامات الآن عبر HTTP interface، لذا تظهر بالقيمة interface = 2 بدلًا من 1، ويحمل الـ column http_user_agent في system.query_log إصدارات Airflow والـ provider إضافةً إلى حقل client_name الإضافي إن كان مضبوطًا. وأي مراقبة كانت تُرشّح بناءً على البروتوكول الأصلي أو على اسم عميل clickhouse-driver تحتاج إلى تحديث. كما أن استعلام SELECT الذي لا يُرجع صفوفًا يُنتج مدخلًا ثانيًا: إذ يُنفّذ مؤشر DB-API العبارة SELECT * FROM (...) LIMIT 0 لاسترجاع البيانات الوصفية للـ column. المُهل الزمنية عبر HTTP. أصبح send_receive_timeout الآن هو مهلة قراءة HTTP، وأي proxy أو موازن حمل بين الـ worker وClickHouse يطبّق مهلة السكون الخاصة به على الطلب. لذا قد تحتاج العبارات التي كانت تُنفَّذ لدقائق طويلة عبر البروتوكول الأصلي إلى رفع تلك الحدود. التعامل مع الاتصالات. يُنشئ الخطاف عميل clickhouse-connect لكل استدعاء run أو get_client، على غرار الطريقة التي كان الـ plugin يفتح بها اتصالًا أصليًا جديدًا لكل execute. ويتشارك العملاء مجمّع اتصالات HTTP على مستوى العملية بأكملها، لذا فإن استدعاء close() على عميل ناتج عن get_client() يُعدّ ممارسة جيدة لا شرطًا إلزاميًا؛ فالمجمّع يُحرَّر عند انتهاء عملية المهمة. والعميل هو context manager، لذا تبقى with hook.get_client() as client: أنظف صيغة.

قائمة التحقق

  1. إصدار Airflow هو 2.11 أو أحدث.
  2. منفذ HTTP يمكن الوصول إليه من الـ workers؛ وشهادات TLS صالحة لنقطة نهاية HTTP.
  3. كل اتصال ClickHouse: النوع clickhouse، المنفذ 8123 أو 8443، والحقول الإضافية مُحوَّلة وفق الخطوة 2، ومُتحقَّق منها باستخدام airflow connections test.
  4. استبدال الـ imports وفق الخطوة 3؛ وحذف بادئات ClickHouse من wrappers الخاصة بـ common.sql.
  5. إعادة تسمية clickhouse_conn_id إلى conn_id في الـ operators والـ sensors، وتعيين conn_id في كل task كان يعتمد على القيمة الافتراضية للـ plugin.
  6. نقل settings= إلى داخل hook_params={"session_settings": ...}.
  7. استبدال hook.execute("INSERT ... VALUES", rows) بـ bulk_insert_rows؛ واستبدال ClickHouseOperator(parameters=rows) بـ SQLInsertRowsOperator.
  8. تعديل الدوال القابلة للنداء في الـ sensors من result[0][0] إلى قيمة الـ cell المجردة.
  9. إعادة كتابة استخدامات with_column_types وexternal_tables وcolumnar وquery_id وtypes_check باستخدام handler أو get_client().
  10. مراجعة الجهات المستهلكة اللاحقة لـ XComs الناتجة عن الـ tasks متعددة العبارات وtasks الإدخال (INSERT).
  11. تحديث الشيفرة التي تعترض استثناءات clickhouse_driver.
  12. إلغاء تثبيت airflow-clickhouse-plugin وclickhouse-driver.

استخدام مساعد برمجي يعمل بالذكاء الاصطناعي

صُمِّم الـ mapping أعلاه ليكون آليًا بحتًا بشكل متعمّد، بحيث يستطيع المساعد البرمجي تطبيقه على repository خاص بـ DAG. وفيما يلي prompt أثبت فعاليته:
راجع الفروقات؛ فتغييرات الاتصال ودلالات المستشعر ومستهلكو XCom هي المواضع التي تخطئ فيها عمليات إعادة الكتابة الآلية.
آخر تعديل في ٢٦ سبتمبر ٢٠٢٦