पाठ 15 / 25
Partitions: repartition, coalesce और Skew
Partitions की संख्या नियंत्रित करें और data skew सँभालें।
सही आकार की कार्य-इकाइयाँ
Partition समानांतरता की इकाई है: प्रति partition एक task। बहुत कम partitions cores को निष्क्रिय छोड़ते हैं और out-of-memory का जोखिम देते हैं; बहुत ज़्यादा छोटे tasks और scheduling का अतिरिक्त बोझ बनाते हैं। आम लक्ष्य है लगभग 100 से 200 MB के partitions और cores की संख्या के कम से कम कुछ गुना। repartition(n) ठीक n संतुलित partitions के लिए पूरा shuffle करता है (और संख्या बढ़ा सकता है); coalesce(n) बिना पूरे shuffle के मौजूदा partitions को सिर्फ़ मिलाता है और सिर्फ़ घटा सकता है। Shuffle के बाद संख्या spark.sql.shuffle.partitions (डिफ़ॉल्ट 200; AQE उसे कम कर सकता है) से आती है। Data skew का मतलब किसी एक key की rows दूसरों से कहीं ज़्यादा हैं, इसलिए एक task लंबा चलता है जबकि बाक़ी इंतज़ार करते हैं; उपायों में AQE skew handling, "गर्म" key की salting, या छोटी ओर को broadcast करना शामिल है।
Partitions गिनना, चलाकर
मैंने यह Apache Spark 4.0.0 (PySpark, local mode, official Docker image में) पर चलाया। local[2] के साथ छोटा DataFrame 2 partitions से शुरू होता है। repartition(8) 8 बनाता है और coalesce(1) 1 तक मिला देता है।
print("default", orders.rdd.getNumPartitions(), "repartition(8)", orders.repartition(8).rdd.getNumPartitions(), "coalesce(1)", orders.repartition(8).coalesce(1).rdd.getNumPartitions())
Output:
default 2 repartition(8) 8 coalesce(1) 1
गर्म key की salting, चलाकर
मैंने यह Apache Spark 4.0.0 (PySpark, local mode, official Docker image में) पर चलाया। Key hot की 8 rows हैं और cold की 1। Random salt (0 से 2) जोड़ने से hot तीन तक उप-समूहों में बँटता है जिन्हें पहले जोड़ा जाता और फिर मिलाया जाता है: योग सही रहते हैं (28 और 1) जबकि काम बँट जाता है।
skew = spark.createDataFrame([("hot", i) for i in range(8)] + [("cold", 1)], ["k", "v"])
salted = skew.withColumn("salt", F.floor(F.rand(seed=1) * 3)).groupBy("k", "salt").agg(F.sum("v").alias("s")).groupBy("k").agg(F.sum("s").alias("total")).orderBy("k")
print([tuple(r) for r in salted.collect()])
Output:
[('cold', 1), ('hot', 28)]त्वरित जाँच: कौन-सी call पूरा shuffle किए बिना partitions घटाती है?
- coalesce(n)
- repartition(n)
- groupBy()
- distinct()
Answer
coalesce(n) — coalesce पड़ोसी partitions को स्थानीय रूप से मिलाता है, जबकि repartition सारा डेटा फिर shuffle करता है।