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

الوحدة 3 — التحويلات والإجراءات والتقييم الكسول

في الوحدة السابقة رأينا كيف يعيد محسِّن Catalyst كتابة استعلاماتنا. لكنّه لا يستطيع فعل ذلك إلّا لأنّ Spark لا يُنفِّذ الشيفرة عند كتابتها. هذه الوحدة تشرح هذا التأجيل، وتُقدِّم أدوات قراءته وأخطاءه الشائعة.

التحويلة والإجراء

كلّ عمليّة على DataFrame تنتمي إلى صنفين:

  • التحويلة (transformation): تُعيد DataFrame جديدًا، دون أن يُحسَب شيء. أمثلة: select، filter، withColumn، groupBy، join.
  • الإجراء (action): يُطلق التنفيذ فعلًا ويُعيد قيمة إلى المحرّك أو يكتب على القرص. أمثلة: count، collect، show، take، write.parquet.

خذ هذا المقطع من مشروعنا:

vols = spark.read.parquet("s3a://donnees-vols/")
retards = vols.filter("retard_arrivee > 15") \
.select("compagnie", "mois", "retard_arrivee") \
.groupBy("compagnie") \
.avg("retard_arrivee")

حتى الآن، لم يقرأ Spark بايتًا واحدًا من Parquet. بنى فقط شجرة عمليّات منطقيّة تنتظر إجراءً. عند retards.show() تنفّذ الشجرة كلّها. عند retards.count() تُنفَّذ مرّة أخرى — نعم، من الصفر.

قراءة خطّة التنفيذ

كلّ DataFrame يعرض explain() الذي يُظهِر ما سيفعله Spark فعلًا:

retards.explain(mode="formatted")

سيُظهِر الخرج أربع خطط: منطقيّة (Parsed Logical Planمُحلَّلة (Analyzedمُحسَّنة (Optimized)، ثم الفيزيائيّة (Physical Plan). الأخيرة هي ما سيُنفَّذ. سترى PushedFilters تُبيِّن أنّ Catalyst دفع مُرشِّحك إلى قارئ Parquet، و ReadSchema تحدّد الأعمدة التي ستُقرأ. من دون قراءة explain تعمل بالخبطة؛ ومَن يُتقن قراءته يوفّر ساعات تصحيح.

تحويلات ضيّقة وأخرى واسعة

ينقسم كلّ تحويل إلى نمطين حسب اتّصال الأقسام:

  • التحويلة الضيّقة: كلّ قسم مُخرَج يعتمد على قسم مُدخَل واحد. أمثلة: filter، select، withColumn. لا حاجة إلى مراسلة بين المنفّذين.
  • التحويلة الواسعة: كلّ قسم مُخرَج يحتاج بيانات من عدّة أقسام مُدخَلة. أمثلة: groupBy، join، orderBy، distinct. يحتاج Spark أن ينقل الأسطر بين المنفّذين ليجمع القيم ذات المفتاح نفسه في قسم واحد. هذه هي عمليّة الخلط (shuffle).

الخلط ثمنه الشبكة والقرص: يكتب كلّ منفّذ ملفّات وسيطة، ثمّ يقرأها المنفّذ التالي عبر الشبكة. على 60 مليون رحلة يمكن لخلط واحد أن يكلّف عدّة غيغابايت من النقل. الخلاصة العمليّة: قلِّل عدد التحويلات الواسعة كلّما استطعت، ورتِّبها لتقع بعد المُرشِّحات لا قبلها.

مثال شائع: تجميع ثمّ ترشيح. اكتبه بالعكس:

# سيّئ: يُخلط 60 مليونًا ثم يُرشَّح
vols.groupBy("compagnie").avg("retard_arrivee").filter("avg(retard_arrivee) > 20")

# جيّد: يُرشَّح أوّلًا فيصغر الخلط
vols.filter("retard_arrivee > 0").groupBy("compagnie").avg("retard_arrivee")

Catalyst يفعل هذا التحسين في كثير من الحالات، لكن ليس كلّها. لا تتّكل عليه، اكتب استعلامك بترتيب معقول.

cache و persist

بما أنّ كلّ إجراء يُعيد التنفيذ من الصفر، فإنّ حساب النموذج نفسه مرّتين مكلف. الحلّ هو تخبئة DataFrame:

retards = retards.cache()
retards.count() # الحساب الأوّل: يقرأ ويُخبِّئ
retards.show() # الاستعمال الثاني: من الذاكرة

cache() مرادف لـ persist(StorageLevel.MEMORY_AND_DISK): يحاول الاحتفاظ بالنتيجة في الذاكرة، ويُسرِّبها إلى القرص إن لَزم. لا تُخبِّئ كلّ شيء: التخبئة نفسها تستهلك ذاكرة وتُنقِص المتاح للمنفّذين. اِحفَظ التخبئة لِـ DataFrame يُستعمَل مرّتين أو أكثر في نفس الوظيفة.

count() ليست حرّة

مبتدئ Spark يتعوَّد على df.count() بعد كلّ تحويل ليتأكّد أنّ الأمور تسير. لكنّ كلّ count هو إجراء كامل يُعيد قراءة كلّ شيء إذا لم تكن DataFrame مُخبَّأة. على مجموعة من 60 مليون سطر، هذا خطأ يحوّل ثوانٍ إلى دقائق. اِحصر الإجراءات في مواضع لك حاجة فعليّة إلى نتيجتها.

الخلاصة

  • التحويلات كسولة، لا تنفذ إلّا عند إجراء؛ كلّ إجراء يُعيد التنفيذ من الصفر بلا تخبئة.
  • explain يُظهِر ما ستفعله Spark فعلًا: دفع المُرشِّحات، تقليم الأعمدة، خوارزميّة الانضمام؛ لا تُطلق مهمّة كبيرة قبل قراءتها.
  • التحويلة الواسعة تُطلِق خلطًا يعبر الشبكة والقرص: قلِّلها ورتِّبها بعد المرشِّحات.
  • cache يفيد حين تُستعمَل نفس DataFrame أكثر من مرّة، لا في كلّ مكان.

الوحدة التالية: قراءة البيانات، ولماذا Parquet يتفوّق على CSV، وكيف يُوفِّر التقسيم بعمود قراءة قد تكون أسرع بعشر مرّات.