प्रदर्शन सुधार

PySpark के साथ डेटा क्लीनिंग

Mike Metzger

Data Engineering Consultant

Spark execution plan समझना

voter_df = df.select(df['VOTER NAME']).distinct()
voter_df.explain()
== Physical Plan ==
*(2) HashAggregate(keys=[VOTER NAME#15], functions=[])
+- Exchange hashpartitioning(VOTER NAME#15, 200)
   +- *(1) HashAggregate(keys=[VOTER NAME#15], functions=[])
      +- *(1) FileScan csv [VOTER NAME#15] Batched: false, Format: CSV, Location: 
      InMemoryFileIndex[file:/DallasCouncilVotes.csv.gz], 
      PartitionFilters: [], PushedFilters: [], 
      ReadSchema: struct<VOTER NAME:string>
PySpark के साथ डेटा क्लीनिंग

Shuffling क्या है?

Shuffling का मतलब है किसी टास्क को पूरा करने के लिए डेटा को अलग-अलग वर्कर तक ले जाना।

  • यह उपयोगकर्ता से जटिलता छिपाता है
  • पूरा होने में धीमा हो सकता है
  • कुल थ्रूपुट घटाता है
  • अक्सर ज़रूरी होता है, पर इसे कम रखने की कोशिश करें
PySpark के साथ डेटा क्लीनिंग

Shuffling कैसे घटाएँ?

  • .repartition(num_partitions) का उपयोग सीमित रखें
    • इसकी जगह .coalesce(num_partitions) करें
  • .join() कॉल करते समय सावधानी रखें
  • .broadcast() का उपयोग करें
  • हो सकता है सीमित करने की ज़रूरत न हो
PySpark के साथ डेटा क्लीनिंग

Broadcasting

Broadcasting:

  • हर वर्कर को ऑब्जेक्ट की एक कॉपी देता है
  • नोड्स के बीच अनावश्यक संचार रोकता है
  • .join() ऑपरेशनों को काफी तेज कर सकता है

.broadcast(<DataFrame>) मेथड उपयोग करें

from pyspark.sql.functions import broadcast
combined_df = df_1.join(broadcast(df_2))
PySpark के साथ डेटा क्लीनिंग

अभ्यास करते हैं!

PySpark के साथ डेटा क्लीनिंग

Preparing Video For Download...