PySpark 집계

PySpark 입문

Benjamin Schmidt

Data Engineer

PySpark SQL 집계 개요

  • 일반적인 SQL 집계는 spark.sql()로 사용할 수 있습니다
    # SQL 집계 쿼리
    spark.sql("""
      SELECT Department, SUM(Salary) AS Total_Salary, AVG(Salary) AS Average_Salary
      FROM employees
      GROUP BY Department
    """).show()
    
PySpark 입문

DataFrame과 SQL 연산 결합

# 3000 초과 급여 필터링
filtered_df = df.filter(df.Salary > 3000)

# 필터링된 DataFrame을 뷰로 등록
filtered_df.createOrReplaceTempView("filtered_employees")

# 필터링된 뷰에 SQL로 집계 수행 spark.sql(""" SELECT Department, COUNT(*) AS Employee_Count FROM filtered_employees GROUP BY Department """).show()
PySpark 입문

집계에서의 데이터 타입 처리

# 타입 캐스팅 예시
data = [("HR", "3000"), ("IT", "4000"), ("Finance", "3500")]
columns = ["Department", "Salary"]
df = spark.createDataFrame(data, schema=columns)

# Salary 컬럼을 정수로 변환 df = df.withColumn("Salary", df["Salary"].cast("int")) # 집계 수행 df.groupBy("Department").sum("Salary").show()
PySpark 입문

집계를 위한 RDD

# RDD로 집계 예시
rdd = df.rdd.map(lambda row: (row["Department"], row["Salary"]))

rdd_aggregated = rdd.reduceByKey(lambda x, y: x + y)
print(rdd_aggregated.collect())
PySpark 입문

PySpark 집계 모범 사례

  • 일찍 필터링: 집계 전 데이터 크기 축소
  • 타입 처리: 데이터 정제 및 올바른 타입 보장
  • 전체 데이터 사용 연산 지양: groupBy() 같은 연산 최소화
  • 올바른 인터페이스 선택: 최적화된 DataFrame을 우선 사용
  • 성능 모니터링: explain()으로 실행 계획 확인 및 최적화
PySpark 입문

핵심 정리

  • PySpark SQL 집계: 데이터를 요약하는 SUM(), AVERAGE() 등의 함수
  • DataFrame과 SQL: 두 접근을 결합해 유연한 데이터 처리
  • 데이터 타입 처리: 집계 중 타입 불일치 문제 해결
  • RDD vs DataFrame: 트레이드오프를 이해하고 적절한 도구 선택
PySpark 입문

연습해 봅시다!

PySpark 입문

Preparing Video For Download...