انتقل إلى المحتوى الرئيسي

الوحدة 9 — كتابة النتائج والتكامل مع الأنظمة اللاحقة

نموذجٌ لا يُنتج تنبُّؤات لا قيمة له. هذه الوحدة عن الطرف الثاني للأنبوب: كيف تُكتب المخرجات بشكل مفيد، كيف يُحفَظ الخطّ المدرَّب لِيُعاد تحميله بلا خطأ في التسجيل، وكيف يجدول كلّ ذلك.

كتابة Parquet

نكتب توقُّعاتنا على تاريخ الرحلات بصيغة Parquet مقسَّمة، بنفس المنطق الذي رأيناه في الوحدة 4 لِلقراءة:

predictions = modele.transform(vols_2024)
predictions.select("date", "compagnie", "vol_id", "prediction", "probability") \
.write \
.mode("overwrite") \
.partitionBy("date") \
.parquet("s3a://predictions-vols/annee=2024/")

prediction و probability هما العمودان اللذان يُضيفهما نموذج تصنيف MLlib افتراضيًّا. أعمدة النموذج الداخليّة (rawPrediction، features) لا نُنتجها لِمستخدمي المخرج: يُختار ما يفيد فقط.

أوضاع الكتابة

mode(...) يقرِّر ماذا يفعل Spark إن وجد وجهة الكتابة موجودة:

  • error (الافتراض): يفشل. آمن، ثقيل في الأنابيب الآليّة.
  • overwrite: يمسح الوجهة ويستبدلها. اِنتبه: يمسح كلّ ما في المسار، لا الملفّات التي كنت ستكتبها فقط.
  • append: يُضيف الملفّات الجديدة إلى الموجود. مفيد لِلبيانات المتراكمة، خطر لأنّه لا يمنع التكرار.
  • ignore: لا يكتب شيئًا إن كانت الوجهة موجودة. مفيد لِلعمليّات التي يجب أن تكون idempotent.

overwrite مع partitionBy له سلوك خاصّ في Spark 3 عبر spark.sql.sources.partitionOverwriteMode="dynamic": يستبدل الأقسام التي تُنتِج أسطرًا فقط، لا الأقسام الأخرى. هذا الوضع مطلوب لِلجدولة اليوميّة: يُعاد حساب يوم واحد بلا مسح البقيّة.

spark.conf.set("spark.sql.sources.partitionOverwriteMode", "dynamic")

# يعيد كتابة أقسام 2024-11-14 و 2024-11-15 فقط
predictions.write.mode("overwrite").partitionBy("date") \
.parquet("s3a://predictions-vols/")

حفظ خطّ الإنتاج المُدرَّب

الأهمّ من التنبُّؤات ذاتها هو حفظ خطّ الإنتاج. PipelineModel قابل لِلحفظ التسلسليّ الكامل:

meilleur_modele = modele_cv.bestModel  # PipelineModel

meilleur_modele.write() \
.overwrite() \
.save("s3a://modeles/retard-vols/v3")

يُنشئ Spark مجلَّدًا كاملًا يحوي مخطَّطًا JSON بترتيب المراحل، وحقيبة لكلّ محوّل مع معلماته وقواميسه، وحقيبة أخيرة للنموذج (الأشجار مثلًا) بصيغة Parquet داخليّة. لا تحرِّر هذه الملفّات يدويًّا.

إعادة التحميل لِلاستدلال

في نصّ منفصل (أو مهمّة مجدوَلة)، نعيد التحميل ونُطبِّق:

from pyspark.ml import PipelineModel

modele = PipelineModel.load("s3a://modeles/retard-vols/v3")
vols_hier = spark.read.parquet("s3a://donnees-vols/date=2024-11-14/")
predictions = modele.transform(vols_hier)

نقطة حاسمة: DataFrame المُدخَل يجب أن يحوي كلّ الأعمدة الأصليّة التي كانت في التدريب، بأنواعها نفسها. النموذج يستدعي كلّ مرحلة بالترتيب، وأوّل مرحلة كانت StringIndexer(inputCol="compagnie") — إن لم يكن العمود موجودًا، سيفشل قبل الوصول إلى الخوارزميّة. لذا اعتنِ بمخطَّط بيانات التسجيل، فهو العَقد بين التدريب والاستدلال.

التسجيل الدفعيّ مقابل التدفّق

خطّ الإنتاج المحفوظ يعمل في نمطَين:

  • دفعيّ (Batch): مهمّة يوميّة تعالج البيانات المستجدَّة. transform على DataFrame كامل، نتيجة إلى Parquet.
  • دفقيّ (Structured Streaming): يقرأ من Kafka أو مجلَّد يستقبل ملفّات، يمرِّر كلّ دفعة عبر PipelineModel.transform، ثمّ يكتب إلى وجهة دفقيّة.

المهمّ: نفس النموذج يعمل في النمطَين بلا تعديل. الاختلاف في جدولة الاستقبال لا في الاستدلال ذاته.

الجدولة

المهام الموزَّعة تحتاج جدولة موثوقة. البدائل الشائعة:

  • Apache Airflow: DAG بلغة بايثون، مشغِّل مُخصَّص لِ Spark (SparkSubmitOperator، DatabricksSubmitRunOperator). المعيار الصناعي حاليًّا.
  • Databricks Jobs: إن كان عنقودك على Databricks، الجدولة مدمجة، مع إعادة محاولة وتنبيهات.
  • AWS Step Functions / GCP Cloud Composer: خيارات سحابيّة مُدارة.

لا تجدولها بـ cron مباشر إن أمكن: لا مراقبة، لا إعادة محاولة، لا اعتماديّات بين المهام.

إصدار النموذج

مسار s3a://modeles/retard-vols/v3 ليس مصادفةً. إصدار النماذج ضرورة عمليّة:

  • رقم الإصدار مع تاريخ التدريب: v3-2024-11-14.
  • بيانات ما تدرَّب عليه النموذج: نطاق التواريخ، عدد الأسطر، Git SHA للشيفرة.
  • مقاييس الأداء على التقييم: AUC، precision-recall، عدد الأسطر المقيَّمة.

هذا كلّه يُخزَّن جنبًا إلى الحقيبة، أو في سِجِلّ نماذج مركزيّ (MLflow Model Registry).

overwrite عدوّك يوم ما

كتابة نموذج جديد بـ overwrite تمسح النسخة السابقة. إن اكتشفت بعد ساعة أن النموذج الجديد أسوأ، لن تعود بسهولة. اِكتب دائمًا في مسار جديد يحوي رقم الإصدار، ثمّ حدِّث اختصارًا (symlink، أو مفتاحًا في قاعدة بيانات) يشير إلى «النسخة الحاليّة». هذا يعطيك التراجع في ثانية.

الخلاصة

  • اُكتب مخرجاتك في Parquet مقسَّمة بعمود ذي دلالة (تاريخ)، لا كامل الأعمدة الداخليّة.
  • أوضاع الكتابة الأربعة تتصرَّف بشكل مختلف؛ overwrite مع partitionOverwriteMode="dynamic" هو الوضع العمليّ لِلجدولة اليوميّة.
  • PipelineModel.save يحفظ السلسلة الكاملة (محوّلات ونموذج) في مجلَّد قابل لِإعادة التحميل بـ PipelineModel.load.
  • إصدار النموذج في مسار جديد لكلّ تدريب، مع سِجِلّ يربط الإصدار بالبيانات والمقاييس، وليس بـ overwrite.

الوحدة التالية: تركيب كلّ ما سبق في مشروع كامل يتوقّع تأخُّرات الرحلات، والمقارنة مع نموذج نظير على pandas.