전체 그래프
Spark

Spark

data-engineeringsparkdistributed-computingbig-data

상위: Data Engineering

요약

Apache Spark는 대규모 분산 데이터 처리 프레임워크입니다. In-Memory 연산으로 빠른 성능을 제공하며, RDD, DataFrame, Spark SQL, Structured Streaming을 지원합니다. 배치 처리와 실시간 처리를 모두 수행할 수 있습니다.

핵심 개념

기본 구조

데이터 구조

  • RDD - Resilient Distributed Dataset 개념과 특징
  • DataFrame - 스키마 기반 분산 데이터 컬렉션
  • Spark SQL - Catalyst Optimizer, Dataset API

연산

고급 기능

학습 경로

1단계: 시작하기 (입문)

  1. Introduction - Spark가 무엇이고 왜 사용하는지
  2. Architecture - Spark가 어떻게 동작하는지
  3. Installation - 설치하고 첫 예제 실행

목표: Spark의 기본 개념과 실행 구조 이해

2단계: 핵심 개념 (기본)

  1. RDD - Spark의 기본 데이터 구조 이해
  2. Transformations - Narrow/Wide 구분, Lazy Evaluation 이해
  3. RDD Transformations - map, filter, groupBy 등 변환 연산 학습
  4. RDD Actions - collect, reduce 등 실행 연산 학습

목표: RDD 기반 데이터 처리 능력 습득

3단계: 고수준 API (실무)

  1. DataFrame - 스키마 기반 데이터 처리
  2. Spark SQL - SQL 쿼리와 Catalyst Optimizer
  3. SQL Operations - 실무 데이터 처리 패턴
  4. Data IO - 다양한 파일 형식 읽기/쓰기

목표: DataFrame/SQL 기반 실무 데이터 처리

4단계: 고급 응용 (심화)

  1. Structured Streaming - 실시간 스트림 처리
  2. Spark Integration - Airflow와 통합
  3. 성능 최적화 및 튜닝

목표: 프로덕션 환경 운영 능력

자주 사용하는 패턴

데이터 읽기

# 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")

관련 주제

참고 자료