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

الوحدة 8 — الضبط: الأقسام والذاكرة والخلط

النموذج مُدرَّب، لكنّه يستغرق ست ساعات بدل نصف ساعة. هذا الوقت لا يفسِّره كود سيّئ بل ضبط سيّئ. هذه الوحدة عن الأدوات التي تحوِّل مهمّة كبيرة إلى مهمّة كبيرة سريعة.

عدد الأقسام

كلّ استعلام يمرّ بأقسام. القسم الواحد مَهمّة واحدة على منفّذ. عدد الأقسام يجب أن يكون 2 إلى 4 أضعاف عدد الأنوية الإجماليّ. أقلّ من ذلك: أنوية عاطلة. أكثر بكثير: كلفة جدولة تُبتلع المكسب.

القيمة الأشهر في التحكّم: spark.sql.shuffle.partitions (الافتراض 200). تُطبَّق بعد كلّ خلط. على عنقود بـ 40 نواة، الافتراض 200 معقول. على 400 نواة تحتاج 800 إلى 1600. على حاسوب محمول بـ 8 أنوية، 200 مبالَغة تنتج جدولة ثقيلة، اِنزل بها إلى 32.

spark = (SparkSession.builder
.config("spark.sql.shuffle.partitions", "32")
.getOrCreate())

يجب ضبطها قبل الاستعلامات، لا في المنتصف.

repartition و coalesce

عمليّتان لِتعديل عدد الأقسام يدويًّا، تختلفان اختلافًا جوهريًّا:

  • repartition(n): يُطلق خلطًا كاملًا لِإعادة توزيع البيانات على n قسم. مكلف، لكن يُعطي أقسامًا متوازنة.
  • coalesce(n): يدمج الأقسام الحاليّة بلا خلط، لكن لا يعمل إلّا لِتقليل العدد. رخيص، ينفع بعد filter قاسٍ قلَّص البيانات.

قاعدة عمليّة:

# بعد ترشيح 60 % من الأسطر، لديّ الآن 200 قسم صغير
donnees_recentes = vols.filter("annee >= 2020")

# سيئ: خلط لا داعي له
donnees_recentes = donnees_recentes.repartition(50)

# جيّد: دمج بلا خلط
donnees_recentes = donnees_recentes.coalesce(50)

repartition يفيد حين تنوي فعل شيء يحتاج بحدّ ذاته إلى خلط. coalesce لِلتقليل قبل الكتابة، حتّى لا تُنتج مئات ملفّات صغيرة.

عدم توازن المفاتيح

مشكلة صامتة لكنّها موجعة: groupBy("compagnie") على تاريخ الرحلات. الشركات الأربع الكبرى تحوي 70 % من الأسطر. سيُخلَط كلّ سطر إلى المنفّذ الذي يحمل مفتاحه، فيُنهي المنفّذون الصغار عملهم في ثوانٍ ويجهد الأربعة المسؤولون عن المفاتيح الشعبيّة لدقائق.

تُقرأ الحادثة في Spark UI: مَهمّة (task) واحدة أو اثنتان أطول عشرات المرّات من البقيّة، بينما «الجُرن» (Timeline) يُظهر منفّذين عاطلين.

الحلّ الكلاسيكيّ الملح (salting):

from pyspark.sql import functions as F

vols_sale = vols.withColumn(
"compagnie_salee",
F.concat_ws("_", F.col("compagnie"),
(F.rand() * 8).cast("int"))
)
resultat = vols_sale.groupBy("compagnie_salee").count()
# ثمّ إعادة تجميع من دون الملح

نضيف مِلْحًا عشوائيًّا يُوزِّع كلّ مفتاح على عدّة أقسام، ثمّ نجمع لاحقًا. الحلّ الحديث: AQE (Adaptive Query Execution) في Spark 3، يُفعَّل بـ spark.sql.adaptive.enabled=true، ويُقسِّم المفاتيح غير المتوازنة تلقائيًّا. اِستعمله دائمًا.

broadcast join

عند ضمّ جدولَين، الافتراض sort-merge join الذي يخلط الاثنين. إن كان أحدهما صغيرًا بما يكفي لِلمنفّذ الواحد (أقلّ من نحو 10 ميغابايت افتراضيًّا، spark.sql.autoBroadcastJoinThreshold)، أرسله بالبث إلى كلّ منفّذ:

from pyspark.sql import functions as F

vols_avec_meteo = vols.join(
F.broadcast(meteo),
on=["date", "aeroport"],
how="left"
)

هذا يُلغي خلط الجدول الكبير، ويكون مكسبه دراميًّا: انضمام كان يستغرق نصف ساعة يصير في ثوانٍ.

ذاكرة المنفّذين

كلّ منفّذ يتقاسم ذاكرته بين تنفيذ وتخزين ووحدة نفقات ثابتة (overhead). الأخيرة (spark.executor.memoryOverhead) هي ما تفوته العينُ عادةً: مكتبات JVM الخارجيّة، بايثون في PySpark، مؤشّرات الشبكة. افتراضها 10 % من ذاكرة المنفّذ. حين ترى OutOfMemoryError رغم أنّ ذاكرة التنفيذ تبدو كافية، ارفع memoryOverhead أوّلًا.

قاعدة تجريبيّة:

  • ذاكرة كلّ منفّذ: spark.executor.memory = 8g إلى 32g. أكثر من 32 يعاني من ضغط الجمع القمامي (GC).
  • memoryOverhead: 15 % لِـPySpark، لأنّ عمليّة بايثون منفصلة عن JVM وتحسب هنا.
  • أنوية لِكلّ منفّذ: 4 إلى 5. أكثر يُدخِل تنافسًا داخل JVM.

قراءة Spark UI

http://<driver>:4040 يعرض واجهة تحوز أربع نوافذ حيويّة:

  • Jobs: قائمة الوظائف الناتجة عن الإجراءات. زمن كلّ وظيفة.
  • Stages: تفصيل داخل الوظيفة. رأس المرحلة يظهر أكبر مَهمّة، متوسّطها، أدناها. إن كان أكبرها 20 ضعف متوسّطها، لديك عدم توازن.
  • Storage: DataFrame مُخبَّأة الآن. نسبتها من الذاكرة.
  • SQL: مخطَّط الاستعلام بترتيبه الفيزيائيّ، بحلقات ملوَّنة تبيِّن أين ضاع الوقت.
قبل أي ضبط، قسها

لا تُغيِّر shuffle.partitions عشوائيًّا. اِفتح Spark UI. اقرأ الرقم الحقيقيّ: أيّ مرحلة أطول، ما توزيع مهامّها، أيّ منفّذ متأخّر. ثمّ عدِّل معلمة واحدة، أَعِد التشغيل، لاحِظ الفرق. الضبط بالحدس على Spark هو الطريق الأقصر لِتضييع أسبوع.

الخلاصة

  • عدد الأقسام 2-4 أضعاف عدد الأنوية؛ افتراض 200 لِshuffle.partitions صالح لعنقود متوسّط، لا لحاسوب محمول.
  • repartition يخلط، coalesce يدمج بلا خلط. اِستعمل الثاني بعد filter قاسٍ وقبل الكتابة.
  • عدم توازن المفاتيح يُقرأ من Spark UI؛ الحلّ الحديث AQE، والكلاسيكي الملح (salting).
  • broadcast join يحوّل ضمَّ جدول صغير إلى عمليّة بلا خلط؛ اِستعمله كلّما كان الجدول الأصغر يسع 10 ميغابايت.

الوحدة التالية: كتابة النتائج، حفظ خطّ الإنتاج، وتكامل الأنابيب مع الأنظمة اللاحقة.