PySpark के साथ डेटा क्लीनिंग
Mike Metzger
Data Engineering Consultant
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>
Shuffling का मतलब है किसी टास्क को पूरा करने के लिए डेटा को अलग-अलग वर्कर तक ले जाना।
.repartition(num_partitions) का उपयोग सीमित रखें.coalesce(num_partitions) करें.join() कॉल करते समय सावधानी रखें.broadcast() का उपयोग करेंBroadcasting:
.join() ऑपरेशनों को काफी तेज कर सकता है.broadcast(<DataFrame>) मेथड उपयोग करें
from pyspark.sql.functions import broadcast
combined_df = df_1.join(broadcast(df_2))
PySpark के साथ डेटा क्लीनिंग