성능 향상

PySpark로 데이터 정제하기

Mike Metzger

Data Engineering Consultant

Spark 실행 계획 설명

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)은 작업을 완료하기 위해 데이터를 여러 워커로 이동하는 것을 말합니다.

  • 사용자에게 복잡성을 숨깁니다
  • 완료에 시간이 걸릴 수 있습니다
  • 전체 처리량을 낮춥니다
  • 종종 필요하지만, 최소화해야 합니다
PySpark로 데이터 정제하기

셔플을 줄이는 방법

  • .repartition(num_partitions) 사용을 최소화합니다
    • 대신 .coalesce(num_partitions)를 사용합니다
  • .join() 호출 시 주의합니다
  • .broadcast()를 사용합니다
  • 항상 제한할 필요는 없습니다
PySpark로 데이터 정제하기

브로드캐스트

브로드캐스트(broadcasting):

  • 각 워커에 객체 사본을 제공합니다
  • 노드 간 불필요한 통신을 막습니다
  • .join()을 크게 가속할 수 있습니다

.broadcast(<DataFrame>) 메서드를 사용합니다

from pyspark.sql.functions import broadcast
combined_df = df_1.join(broadcast(df_2))
PySpark로 데이터 정제하기

실습해 봅시다!

PySpark로 데이터 정제하기

Preparing Video For Download...