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

الوحدة 10 — مشروع: نموذج مُدرَّب على مجموعة بيانات كبيرة

عبرنا تسع وحدات نظريّة وعمليّة. الآن نصل إلى الأمانة النهائيّة: هل ما تعلّمناه فعلًا يستحقّ ما دفعناه من تعقيد؟ نبني السلسلة كاملة على تاريخ الرحلات، ونقارنها بنموذج نظير على pandas مع عيّنة، ونحكم بلا مجاملة.

البيانات والهدف

ملفّ Parquet مقسَّم بالسنة والشهر، 60 مليون سطر بين 2015 و 2023، حوالي 25 غيغابايت مضغوطة. الهدف الثنائي en_retard بقيمة 1 إذا كان تأخّر الوصول أكثر من 15 دقيقة. المتغيّرات المتاحة: compagnie، origine، destination، distance، mois، jour_semaine، heure_depart_planifiee، meteo_origine، وأخريات.

قراءة وتقسيم

from pyspark.sql import SparkSession, functions as F

spark = (SparkSession.builder
.appName("VolsRetardV3")
.config("spark.sql.adaptive.enabled", "true")
.config("spark.sql.shuffle.partitions", "400")
.getOrCreate())

vols = spark.read.parquet("s3a://donnees-vols/")

# التقسيم زمنيّ لا عشوائيّ: التدريب على 2015-2022، الاختبار على 2023
train = vols.filter("annee <= 2022")
test = vols.filter("annee == 2023")

print(f"تدريب: {train.count():,} سطر")
print(f"اختبار: {test.count():,} سطر")

التقسيم الزمنيّ لا العشوائيّ ضروريّ: في الإنتاج نتوقّع المستقبل من الماضي، والاختبار العشوائي على البيانات الزمنيّة يُنتج مقاييس متفائلة كذبًا.

بناء الـPipeline

نُغلِّف كلّ خطواتنا وفقًا لِلوحدة 5 و 6:

from pyspark.ml import Pipeline
from pyspark.ml.feature import StringIndexer, OneHotEncoder, VectorAssembler
from pyspark.ml.classification import GBTClassifier

# ترميز التصنيفيّات
indexer_c = StringIndexer(inputCol="compagnie", outputCol="c_idx",
handleInvalid="keep")
indexer_o = StringIndexer(inputCol="origine", outputCol="o_idx",
handleInvalid="keep")
indexer_d = StringIndexer(inputCol="destination", outputCol="d_idx",
handleInvalid="keep")

encoder = OneHotEncoder(
inputCols=["c_idx", "o_idx", "d_idx"],
outputCols=["c_oh", "o_oh", "d_oh"]
)

assembleur = VectorAssembler(
inputCols=["c_oh", "o_oh", "d_oh", "distance", "mois",
"jour_semaine", "heure_depart_planifiee"],
outputCol="features",
handleInvalid="skip"
)

gbt = GBTClassifier(labelCol="en_retard", featuresCol="features",
maxDepth=6, maxIter=100)

pipeline = Pipeline(stages=[indexer_c, indexer_o, indexer_d,
encoder, assembleur, gbt])

ملاحظة: لم نستعمل StandardScaler لأنّ GBT شجريّ ولا يستفيد من التسوية.

الضبط والتدريب

شبكة صغيرة، مطويّتان فقط لِلاحتفاظ بالوقت معقولًا:

from pyspark.ml.tuning import TrainValidationSplit, ParamGridBuilder
from pyspark.ml.evaluation import BinaryClassificationEvaluator

grille = (ParamGridBuilder()
.addGrid(gbt.maxDepth, [5, 7])
.addGrid(gbt.maxIter, [80, 120])
.build())

evaluateur = BinaryClassificationEvaluator(labelCol="en_retard",
metricName="areaUnderROC")

tvs = TrainValidationSplit(
estimator=pipeline,
estimatorParamMaps=grille,
evaluator=evaluateur,
trainRatio=0.8,
parallelism=4,
seed=42
)

modele = tvs.fit(train)

على عنقود من 8 عقد، كلٌّ منها 16 نواة و 64 غيغابايت، هذا يستغرق نحو 40 دقيقة. على Databricks بمعالج آليّ لِلعنقود التكلفة نحو 12 دولارًا.

التقييم

predictions = modele.transform(test)
auc = evaluateur.evaluate(predictions)
print(f"AUC على 2023: {auc:.4f}")

نتيجة نموذجيّة: AUC = 0.712. ليس ممتازًا، لكنّه صادق: تأخُّرات الطيران فيها عشوائيّة كبيرة (طقس، حركة جوّيّة، حوادث).

المقارنة مع pandas

الآن نُقيَّم عالجه. نأخذ عيّنة عشوائيّة من 5 % من التدريب (نحو 2.5 مليون سطر) على آلة واحدة قويّة، ونُدرِّب LightGBM:

import pandas as pd
from lightgbm import LGBMClassifier
from sklearn.model_selection import train_test_split
from sklearn.metrics import roc_auc_score

# نُصدِّر عيّنة إلى Parquet ثمّ نقرؤها بـpandas
train.sample(0.05, seed=42).write.mode("overwrite") \
.parquet("/local/echantillon_vols.parquet")

df = pd.read_parquet("/local/echantillon_vols.parquet")
X = pd.get_dummies(df.drop(columns=["en_retard"]))
y = df["en_retard"]

X_tr, X_val, y_tr, y_val = train_test_split(X, y, test_size=0.2,
random_state=42)

modele_local = LGBMClassifier(max_depth=6, n_estimators=200)
modele_local.fit(X_tr, y_tr)

auc_local = roc_auc_score(y_val, modele_local.predict_proba(X_val)[:, 1])
print(f"AUC LightGBM على عيّنة 5 %: {auc_local:.4f}")

نتيجة نموذجيّة: AUC = 0.702 على العيّنة، 0.687 على 2023 الكامل بعد التقييم النهائيّ. الفرق مع Spark GBT هو نقطتان مئويّتان فقط.

ماذا يقول هذا الفرق؟

Spark ML أضاف نقطتَين على AUC مقابل:

  • عنقود لِعشرة أشخاص بدل حاسوب محمول
  • 40 دقيقة تدريب بدل 8 دقائق
  • 12 دولارًا لكلّ تدريب بدل صفر تكلفة سحابيّة

هذه المفاضلة تُتَّخذ بأرقام العمل، لا بمشاعر تقنيّة:

  • إذا كان النموذج يُوَقِّع قرارًا تسويقيًّا يخصّ 200 مليون دولار مبيعات سنويًّا، النقطتان تستحقّان.
  • إذا كان يقود لوحة معلومات داخليّة تخدم عشرة موظّفين، نموذج pandas على عيّنة كافٍ.

متى فاز الحجم فعلًا

يفوز التوزيع بوضوح في حالة أخرى: عدد المفاتيح. لو كنّا نُدرِّب نموذجًا مخصَّصًا لكلّ زوج (مطار مغادرة، شهر) — نحو 4000 نموذج — فالتدريب على آلة واحدة مستحيل. Spark يوزّع الأربعة آلاف عبر المنفّذين طبيعيًّا. هذا هو المكسب الحقيقيّ لِلتوزيع، أكثر من مجرّد حجم البيانات.

قاعدة القرار النهائيّة

Spark ML ليس ترقية طبيعيّة لِ pandas. هو أداة أخرى لِمسألة أخرى. اِستعمله حين تكون الإجابة على مقياس مشكلتك المستحيلة أو تكون البيانات نفسها موزَّعة أصلًا في بحيرة بيانات ولا يمكن أن تعبرها إلى آلة واحدة. لِلمسائل التي تسع عيّنة صادقة في ذاكرة آلة، ابقَ في pandas + LightGBM: النموذج أقوى والتصحيح أسرع.

الخلاصة

  • السلسلة الكاملة تعمل على 60 مليون سطر: قراءة موزَّعة، تمهيد داخل Pipeline، تدريب GBT، حفظ النموذج.
  • AUC = 0.712 مقابل 0.687 لِ LightGBM على عيّنة 5 %؛ فرق نقطتين بتكلفة عنقود.
  • مبرِّر Spark ML الحقيقيّ: البيانات لا تسع آلة واحدة، أو المسألة تحتاج مئات النماذج بالتوازي.
  • إن سعت البيانات في ذاكرة عيّنة صادقة، ابقَ في pandas + LightGBM؛ Spark ML قرار عمل، لا شهادة تقنيّة.

انتهت الوحدات. الوحدة الحادية عشرة تُلخِّص كلّ ما رأيناه وتُقدِّم للامتحان النهائي.