पाठ 13 / 25
Narrow और Wide Transformations: Shuffle
समझें कि कुछ ऑपरेशन network पर डेटा क्यों ले जाते हैं और सबसे महँगे क्यों हैं।
महँगा चरण
Narrow transformation (filter, select, withColumn, map) हर output partition की गणना एक ही input partition से कर सकता है, इसलिए डेटा हिलाए बिना चलता है। Wide transformation (groupBy, join, distinct, repartition, orderBy, partitionBy वाली window) को संबंधित keys वाली rows का मिलना ज़रूरी है, इसलिए Spark shuffle करता है: बीच का डेटा स्थानीय disk पर लिखता है, network पर भेजता है और वापस पढ़ता है। Shuffle एक stage ख़त्म करता और अगला शुरू करता है। Spark ट्यूनिंग का ज़्यादातर हिस्सा कम और छोटे shuffles करने के बारे में है। Plan में shuffle Exchange operator के रूप में दिखता है।
Plan में shuffle देखना, चलाकर
मैंने यह Apache Spark 4.0.0 (PySpark, local mode, official Docker image में) पर चलाया। Exchange hashpartitioning(city, 4) पंक्ति shuffle है। Spark पहले हर partition पर आंशिक aggregation (partial_sum) करता है, city से 4 partitions में shuffle करता है, फिर अंतिम aggregation करता है। # संख्याएँ आंतरिक IDs हैं और runs के बीच बदलती हैं।
orders.groupBy("city").agg(F.sum("amount")).explain()
Output:
== Physical Plan ==
AdaptiveSparkPlan isFinalPlan=false
+- HashAggregate(keys=[city#2], functions=[sum(amount#3)])
+- Exchange hashpartitioning(city#2, 4), ENSURE_REQUIREMENTS, [plan_id=66]
+- HashAggregate(keys=[city#2], functions=[partial_sum(amount#3)])
+- Project [city#2, amount#3]
+- Scan ExistingRDD[order_id#0L,customer#1,city#2,amount#3,order_date#4]Wide ऑपरेशनों से पहले filter और select करें
Shuffle में कम rows और columns जाने का मतलब कम डेटा लिखा और ले जाया जाना। Catalyst अक्सर filters को आपके लिए नीचे ले जाता है, पर जल्दी column चयन और filtering स्पष्ट रूप से करें।
त्वरित जाँच: कौन-सा ऑपरेशन आम तौर पर shuffle कराता है?
- दो columns का select
- groupBy aggregation
- किसी column पर filter
- स्थिरांक के साथ withColumn
Answer
groupBy aggregation — Grouping को किसी key की सारी rows साथ चाहिए, इसलिए डेटा partitions के बीच जाना ज़रूरी है।