الوحدة 6 — خطوط المعالجة والتحقّق المتقاطع الموزَّع
في الوحدة السابقة أَنتَجْنا محوّلاتنا وقلنا سنغلِّفها في Pipeline. الآن نبني السلسلة، ونُشغِّل ضبطًا موزَّعًا على شبكة معاملات، ونتفادى الفخّ الأخطر في MLlib: تسرُّب بيانات التحقّق إلى التمهيد.
لماذا نحتاج Pipeline
المشكلة العمليّة: كلّ محوّل يتعلّم شيئًا من التدريب (قاموس، متوسّط، انحراف). إن طبَّقنا التمهيد على التدريب والتحقّق مجتمعين، تسرَّبت إحصاءات التحقّق إلى المتوسّطات والانحرافات. النموذج يرى معلومات لن يراها في الإنتاج. وحين نُشغِّل تحقّقًا متقاطعًا فوق هذا، تصير كلّ طيّة (fold) ملوَّثة قليلًا، وتصير مقاييس الأداء متفائلة بلا سبب.
Pipeline يحلّ هذا: يجمع مقدِّراتٍ ومحوّلات في مقدِّر واحد. حين تناديه بـ fit(df_entrainement) يُطبِّق fit على كلّ خطوة بترتيبها، ثمّ يُنتج PipelineModel يمكن تطبيقه بـ transform على أيّ DataFrame جديدة. القاعدة الذهبيّة: كلّ خطوة معالجة تتعلّم شيئًا من البيانات يجب أن تعيش داخل الـPipeline.
from pyspark.ml import Pipeline
from pyspark.ml.feature import StringIndexer, OneHotEncoder, VectorAssembler, StandardScaler
from pyspark.ml.classification import GBTClassifier
indexer = StringIndexer(inputCol="compagnie", outputCol="compagnie_idx",
handleInvalid="keep")
encoder = OneHotEncoder(inputCol="compagnie_idx", outputCol="compagnie_oh")
assembleur = VectorAssembler(
inputCols=["compagnie_oh", "distance", "mois", "heure_depart"],
outputCol="features_brutes")
scaler = StandardScaler(inputCol="features_brutes", outputCol="features",
withMean=True, withStd=True)
gbt = GBTClassifier(labelCol="en_retard", featuresCol="features")
pipeline = Pipeline(stages=[indexer, encoder, assembleur, scaler, gbt])
خمس خطوات، مقدِّر واحد. pipeline.fit(train) يُنفّذ كلّ شيء بالترتيب.
شبكة المعاملات
ParamGridBuilder يبني الشبكة التي سنستكشفها:
from pyspark.ml.tuning import ParamGridBuilder
grille = (ParamGridBuilder()
.addGrid(gbt.maxDepth, [4, 6, 8])
.addGrid(gbt.maxIter, [50, 100])
.addGrid(gbt.stepSize, [0.05, 0.1])
.build())
عدد المجموعات: 3 × 2 × 2 = 12 مجموعة. المهمّ التنبُّه: هذا ليس عدد النماذج المُدرَّبة، بل يُضرَب في عدد الطيّات. مع خمس طيّات تصير 60 نموذجًا كاملًا يُدرَّب كلّ منها على 48 مليون سطر. اجعل الشبكة صغيرة أوّلًا، ثمّ وسِّعها.
CrossValidator
يستقبل الـPipeline والشبكة ومقيّمًا:
from pyspark.ml.tuning import CrossValidator
from pyspark.ml.evaluation import BinaryClassificationEvaluator
evaluateur = BinaryClassificationEvaluator(
labelCol="en_retard",
metricName="areaUnderROC"
)
cv = CrossValidator(
estimator=pipeline,
estimatorParamMaps=grille,
evaluator=evaluateur,
numFolds=5,
parallelism=4,
seed=42,
)
modele_cv = cv.fit(train)
parallelism=4 يُخبِر CrossValidator بتدريب أربع مجموعات معاملات في وقت واحد. Spark يوزّع كلّ تدريب على المنفّذين على أيّ حال، لكن مع هذه المعلمة يستهلك أربعة تدريبات نظريًّا على أربع مجموعات مستقلّة من الأنوية. اضبط هذا حسب حجم عنقودك: parallelism=1 يقتل المكسب، وقيمة مبالَغة تحدث ازدحامًا. أربعة معقولة لعنقود متوسّط.
TrainValidationSplit: البديل السريع
خمس طيّات ضربًا في 12 معلمة ضربًا في 48 مليون سطر مكلف. أحيانًا تكفي طيّة واحدة كبيرة (TrainValidationSplit بنسبة 0.8/0.2). يُوَصَّى به:
- في مراحل الاستكشاف قبل التنقية النهائيّة
- حين تكون البيانات ضخمة جدًّا (فوق المئة مليون سطر)
- حين يكون النموذج نفسه مكلفًا (شبكات عصبيّة، Gradient Boosting بأعماق كبيرة)
كلّ ما هو باقٍ من المنطق نفسه، والواجهة مطابقة تقريبًا.
تسرُّب البيانات: الخطأ الذي لا يُغتَفَر
الفخّ الأشهر هو التمهيد خارج الـPipeline:
# سيّئ جدًّا: تسرُّب مضمون
scaler = StandardScaler(...)
data_normalisee = scaler.fit(vols_complet).transform(vols_complet)
train, test = data_normalisee.randomSplit([0.8, 0.2])
cv.fit(train) # المتوسّط والانحراف حُسبا على test أيضًا!
المتوسّط والانحراف حُسبا على vols_complet الذي يشمل بيانات الاختبار. النموذج تسرَّبت له معلومات لن يراها في الإنتاج. المقاييس تكذب. لا يوجد أيّ إنذار.
الطريقة الصحيحة:
# جيّد: التمهيد جزء من الـPipeline
train, test = vols.randomSplit([0.8, 0.2], seed=42)
modele_cv = cv.fit(train) # التمهيد يتعلّم من train فقط
predictions = modele_cv.transform(test)
هذا يشمل كلّ خطوة تتعلّم: StringIndexer (القاموس)، StandardScaler (المتوسّط)، Imputer (الوسيط)، وأيّ استعمال يدويّ لـ .mean() أو .groupBy().avg() قبل التقسيم.
modele_cv.bestModel هو PipelineModel جاهز للاستدلال، يحوي جميع المحوّلات و GBT بأفضل معاملات. اِحفظه بـ modele_cv.bestModel.write().overwrite().save(...) لتعيد تحميله لاحقًا. تخزينه يشمل كلّ ما يلزم: قواميس المُفَهرِسات، المتوسّطات، والأشجار.
قاعدة عمليّة قبل الانتقال
قبل أن تُطلق أيّ CrossValidator، ضع لنفسك ثلاثة حواجز صريحة: أوّلًا، احسب حجم الشبكة يدويًّا (مجموعات × طيّات × parallelism) وقارن الناتج بميزانيّتك الزمنيّة؛ ثانيًا، تأكّد من أنّ كلّ خطوة تعلّم موجودة داخل الـPipeline لا خارجه، حتّى Imputer الذي يبدو بريئًا؛ ثالثًا، اِحفظ bestModel فور انتهاء التدريب في مسار موسوم بالتاريخ لتتمكّن من العودة إليه بعد أسبوعين دون أن تُعيد كلّ التجربة. هذه الحواجز الثلاثة تُوفّر ساعات ضائعة في تشخيص نتائج مشبوهة.
الخلاصة
Pipelineيجمع محوّلات ومقدِّر نهائيّ في مقدِّر واحد، يضمن أن كلّ خطوة تتعلّم من التدريب فقط.ParamGridBuilderيبني شبكة المعاملات؛ عدد التدريبات = مجموعات × طيّات، فاِبدأ صغيرًا.CrossValidatorينفّذ التحقّق المتقاطع الموزَّع؛parallelismيتحكّم في عدد الأنابيب المتوازية.- تسرُّب البيانات ينتج حين يخرج التمهيد من الـPipeline؛ لا تنتج شيئًا من الإحصاءات قبل التقسيم.
الوحدة التالية: ماذا نجد فعلًا في MLlib، وماذا ل ا نجده، وكيف يُعوَّض الغائب حين تحتاجه.