대규모 PySpark

PySpark 입문

Benjamin Schmidt

Data Engineer

스케일 활용

  • PySpark는 GB~TB 규모 데이터에 효과적입니다.
  • 목표: 빠르고 효율적인 처리
  • 실행 방식 이해로 추가 효율 확보
  • 전체 클러스터 관리를 위해 브로드캐스트 사용
joined_df = large_df.join(broadcast(small_df), 
                          on="key_column", how="inner")
joined_df.show()
PySpark 입문

실행 계획

# Using explain() to view the execution plan
df.filter(df.Age > 40).select("Name").explain()
== Physical Plan ==
*(1) Filter (isnotnull(Age) AND (Age > 30))
+- Scan ExistingRDD[Name:String, Age:Int]
1 https://spark.apache.org/docs/latest/api/python/reference/pyspark.sql/api/pyspark.sql.DataFrame.explain.html
PySpark 입문

DataFrame 캐싱과 퍼시스트

  • 캐싱: 작은 데이터셋을 메모리에 저장해 빠르게 접근
  • 퍼시스트: 더 큰 데이터셋을 다양한 저장 레벨로 저장
df = spark.read.csv("large_dataset.csv", header=True, inferSchema=True)

# Cache the DataFrame
df.cache()

# Perform multiple operations on the cached DataFrame df.filter(df["column1"] > 50).show() df.groupBy("column2").count().show()
PySpark 입문

다양한 저장 레벨로 퍼시스트

# Persist the DataFrame with storage level
from pyspark import StorageLevel

df.persist(StorageLevel.MEMORY_AND_DISK)

# Perform transformations result = df.groupBy("column3").agg({"column4": "sum"}) result.show() # Unpersist after use df.unpersist()
PySpark 입문

PySpark 최적화

  • 소규모 하위 구간: 데이터가 많을수록 느려집니다. 선택성이 높은 map()groupby()보다 우선 사용
  • 브로드캐스트 조인: 작은 데이터여도 클러스터 리소스를 모두 사용함
  • 반복 액션 회피: 동일 데이터에 대한 반복 액션은 시간·리소스 낭비
PySpark 입문

Passons à la pratique !

PySpark 입문

Preparing Video For Download...