← 목록으로
단계 07

Spark 입문 · 처리 · 심화

Spark 등장 배경, RDD vs DataFrame, PySpark 설치와 첫 실행, DataFrame·Spark SQL, 파티셔닝·캐싱 등 성능 개념.

왜 Spark인가

MapReduce는 단계마다 디스크에 중간 결과를 씁니다. Spark는 메모리에 두고 여러 단계를 이어서 처리해 10~100배 빠릅니다. 또한 SQL·스트리밍·ML을 하나의 엔진에서 다룹니다.

설치 & 첫 실행

pip install pyspark
from pyspark.sql import SparkSession

spark = SparkSession.builder.appName("intro").master("local[*]").getOrCreate()
df = spark.read.csv("sales.csv", header=True, inferSchema=True)
df.printSchema()
df.show(5)

local[*]: 내 컴퓨터의 모든 코어로 실행. 클러스터에선 yarn 이나 spark://....

RDD vs DataFrame

RDDDataFrame
형태객체의 분산 컬렉션이름 있는 열을 가진 테이블
최적화없음 (내가 쓴 대로)Catalyst 옵티마이저가 자동 최적화
언제저수준 제어가 필요할 때거의 항상
# RDD — MapReduce와 똑같은 사고
rdd = spark.sparkContext.textFile("text.txt")
counts = rdd.flatMap(lambda l: l.split()).map(lambda w: (w, 1)).reduceByKey(lambda a, b: a + b)

# DataFrame — 같은 일
from pyspark.sql import functions as F
words = spark.read.text("text.txt").select(F.explode(F.split("value", " ")).alias("w"))
words.groupBy("w").count().show()

핵심 개념: Transformation과 Action, 지연 실행

  • Transformation (filter, select, groupBy, join): 계획만 세움. 실행 안 함.
  • Action (show, count, collect, write): 이때 비로소 실행.
  • 덕분에 Spark는 전체 계획을 보고 최적화한 뒤 한 번에 실행합니다. df.explain()으로 계획 확인.

DataFrame 처리 (9주차)

from pyspark.sql import functions as F

df = spark.read.parquet("sales.parquet")

# 필터 · 선택 · 새 열
df2 = df.filter(F.col("amount") > 100).select("region", "amount").withColumn("krw", F.col("amount") * 1350)

# 집계
df.groupBy("region").agg(F.sum("amount").alias("total"), F.count("*").alias("n")).orderBy(F.desc("total")).show()

# 조인
regions = spark.read.csv("regions.csv", header=True)
df.join(regions, on="region", how="left").show()

# Spark SQL — 같은 일을 SQL로
df.createOrReplaceTempView("sales")
spark.sql("SELECT region, SUM(amount) total FROM sales GROUP BY region ORDER BY total DESC").show()

# 저장
df2.write.mode("overwrite").partitionBy("region").parquet("out/sales")

심화: 성능 개념 (10주차)

파티션

DataFrame은 여러 파티션으로 나뉘어 각 executor가 병렬 처리합니다.

df.rdd.getNumPartitions()
df.repartition(8)      # 셔플하며 재분배 (늘리거나 고르게)
df.coalesce(1)         # 셔플 없이 줄이기 (작은 결과 저장 시)
  • 파티션이 너무 적으면 병렬성 부족, 너무 많으면 작은 태스크 오버헤드
  • 저장 시 partitionBy("date") → 읽을 때 날짜 조건이 있으면 그 폴더만 읽음 (partition pruning)

셔플 (Shuffle)

groupBy, join, repartition같은 키를 같은 노드로 모으는 네트워크 이동(셔플) 이 발생 — 가장 비싼 작업. 가능한 한 줄이는 게 튜닝의 핵심.

  • 작은 테이블 조인은 broadcast join: df.join(F.broadcast(small), "k")
  • 필터를 먼저, 조인은 나중에 (Catalyst가 대부분 알아서 하지만 확인)

캐싱

같은 DataFrame을 여러 action에 쓰면 매번 다시 계산됩니다.

df_clean = df.filter(...).select(...)
df_clean.cache()           # 또는 .persist()
df_clean.count()           # 이때 메모리에 올라감
df_clean.groupBy("a").count().show()   # 재계산 없이 빠름
df_clean.unpersist()

대용량 처리 체크리스트

  1. 입력은 Parquet, 필요한 열만 select
  2. 필터를 최대한 앞에
  3. 셔플 횟수 확인 (explain()에서 Exchange)
  4. 스큐(특정 키에 데이터 쏠림) 확인 — Spark UI(localhost:4040) 태스크 시간 편차
  5. 반복 사용하는 중간 결과는 cache()

연습 과제

  1. Pandas로 했던 집계를 PySpark DataFrame과 Spark SQL 두 방식으로 재현
  2. explain() 결과에서 Exchange(셔플)가 몇 번 나오는지 세어보기
  3. cache() 전후로 두 번째 action의 시간 비교