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로 데이터 정제하기