상위: Data Engineering
요약
Apache Spark는 대규모 분산 데이터 처리 프레임워크입니다. In-Memory 연산으로 빠른 성능을 제공하며, RDD, DataFrame, Spark SQL, Structured Streaming을 지원합니다. 배치 처리와 실시간 처리를 모두 수행할 수 있습니다.
핵심 개념
기본 구조
- Introduction - Spark 개요, 특징, 탄생 배경
- Architecture - Driver, Executor, Cluster Manager, Job, Stage, Task
- Installation - 설치 및 환경 설정
데이터 구조
- RDD - Resilient Distributed Dataset 개념과 특징
- DataFrame - 스키마 기반 분산 데이터 컬렉션
- Spark SQL - Catalyst Optimizer, Dataset API
연산
- Transformations - Narrow vs Wide, Lazy Evaluation
- RDD Transformations - map, filter, groupBy, join 등
- RDD Actions - collect, reduce, count, sum 등
- SQL Operations - SELECT, GROUP BY, JOIN, FILTER 등
고급 기능
- Structured Streaming - Kafka 통합, 실시간 처리
- Data IO - 파일 읽기/쓰기, 다양한 데이터 소스
- Spark Integration - Airflow와 Spark 연동
학습 경로
1단계: 시작하기 (입문)
- Introduction - Spark가 무엇이고 왜 사용하는지
- Architecture - Spark가 어떻게 동작하는지
- Installation - 설치하고 첫 예제 실행
목표: Spark의 기본 개념과 실행 구조 이해
2단계: 핵심 개념 (기본)
- RDD - Spark의 기본 데이터 구조 이해
- Transformations - Narrow/Wide 구분, Lazy Evaluation 이해
- RDD Transformations - map, filter, groupBy 등 변환 연산 학습
- RDD Actions - collect, reduce 등 실행 연산 학습
목표: RDD 기반 데이터 처리 능력 습득
3단계: 고수준 API (실무)
- DataFrame - 스키마 기반 데이터 처리
- Spark SQL - SQL 쿼리와 Catalyst Optimizer
- SQL Operations - 실무 데이터 처리 패턴
- Data IO - 다양한 파일 형식 읽기/쓰기
목표: DataFrame/SQL 기반 실무 데이터 처리
4단계: 고급 응용 (심화)
- Structured Streaming - 실시간 스트림 처리
- Spark Integration - Airflow와 통합
- 성능 최적화 및 튜닝
목표: 프로덕션 환경 운영 능력
자주 사용하는 패턴
데이터 읽기
# CSV 읽기
df = spark.read.csv("data.csv", header=True, inferSchema=True)
# Parquet 읽기 (추천)
df = spark.read.parquet("data.parquet")
# JSON 읽기
df = spark.read.json("data.json")
데이터 변환
# 필터링 + 선택
result = df.filter(col("age") > 30).select("name", "age")
# 집계
agg_result = df.groupBy("category").agg(
count("*").alias("count"),
avg("price").alias("avg_price")
)
# 조인
joined = df1.join(df2, "key", "inner")
데이터 저장
# Parquet 저장 (추천)
df.write.mode("overwrite").parquet("output.parquet")
# 파티셔닝
df.write.partitionBy("year", "month").parquet("output.parquet")
스트리밍 처리
# Kafka 읽기
streaming_df = spark.readStream \
.format("kafka") \
.option("kafka.bootstrap.servers", "localhost:9092") \
.option("subscribe", "topic") \
.load()
# 집계 및 저장
query = streaming_df \
.groupBy(window("timestamp", "1 minute")) \
.count() \
.writeStream \
.outputMode("update") \
.format("console") \
.start()
주의사항 / 함정
1. collect() 사용 주의
# ❌ Bad: 대용량 데이터에서 메모리 부족
result = df.collect()
# ✅ Good: 샘플링 또는 파일로 저장
sample = df.take(10)
df.write.parquet("output.parquet")
2. Shuffle 최소화
# ❌ Bad: 불필요한 shuffle
df.repartition(100).filter(...)
# ✅ Good: 필터링 후 repartition
df.filter(...).repartition(10)
3. 스키마 명시
# ❌ Bad: inferSchema는 전체 파일 스캔
df = spark.read.csv("data.csv", inferSchema=True)
# ✅ Good: 스키마 명시
df = spark.read.csv("data.csv", schema=predefined_schema)
4. GroupByKey vs ReduceByKey
# ❌ Bad: groupByKey는 모든 데이터 shuffle
rdd.groupByKey().mapValues(sum)
# ✅ Good: reduceByKey는 로컬 집계 후 shuffle
rdd.reduceByKey(lambda a, b: a + b)
5. 캐싱 활용
# ❌ Bad: 같은 연산 반복
result1 = df.filter(...).count()
result2 = df.filter(...).sum()
# ✅ Good: 중간 결과 캐싱
cached = df.filter(...).cache()
result1 = cached.count()
result2 = cached.sum()
6. 파티션 관리
# 파일 크기 128MB ~ 1GB 권장
# 너무 작은 파일: 메타데이터 오버헤드
# 너무 큰 파일: 메모리 부족
# 적절한 파티션 수 설정
df.repartition(200).write.parquet("output")
도구 비교
Spark vs Flink
- Spark: 마이크로 배치, ML/배치 통합, 낮은 진입 장벽
- Flink: 진짜 실시간, 초저지연, 복잡한 이벤트 처리
→ 자세한 비교는 Flink vs Spark 참고
Spark vs Hadoop MapReduce
- Spark: In-Memory, 빠름, 통합 API
- Hadoop: 디스크 기반, 느림, 안정성
활용 사례
데이터 웨어하우스 적재 (ETL)
# CSV → Parquet 변환 및 적재
df = spark.read.csv("raw_data.csv", header=True)
df.write.mode("overwrite").parquet("s3://warehouse/processed/")
로그 데이터 분석
# 로그 집계 및 분석
logs = spark.read.text("logs/*.log")
error_logs = logs.filter(col("value").contains("ERROR"))
error_logs.groupBy(date_format("timestamp", "yyyy-MM-dd")).count()
머신러닝 데이터 전처리
# 대량 데이터 전처리
from pyspark.ml.feature import VectorAssembler
assembler = VectorAssembler(inputCols=features, outputCol="features")
train_data = assembler.transform(df)
실시간 추천 시스템
# Kafka 스트림 처리
user_events = spark.readStream.format("kafka").load()
recommendations = user_events.join(item_features, "item_id")
관련 주제
- Airflow - Spark 작업 오케스트레이션
- Spark Integration - Airflow-Spark 연동 패턴
- Flink - 실시간 처리 대안
- Batch vs Real-time - 처리 방식 비교