← 목록으로
단계 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
| RDD | DataFrame | |
|---|---|---|
| 형태 | 객체의 분산 컬렉션 | 이름 있는 열을 가진 테이블 |
| 최적화 | 없음 (내가 쓴 대로) | 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()
대용량 처리 체크리스트
- 입력은 Parquet, 필요한 열만
select - 필터를 최대한 앞에
- 셔플 횟수 확인 (
explain()에서Exchange) - 스큐(특정 키에 데이터 쏠림) 확인 — Spark UI(
localhost:4040) 태스크 시간 편차 - 반복 사용하는 중간 결과는
cache()
연습 과제
- Pandas로 했던 집계를 PySpark DataFrame과 Spark SQL 두 방식으로 재현
explain()결과에서Exchange(셔플)가 몇 번 나오는지 세어보기cache()전후로 두 번째 action의 시간 비교