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