상위: Spark
요약
Spark SQL은 구조화된 데이터를 처리하기 위한 모듈로, DataFrame의 핵심 엔진입니다. Catalyst Optimizer를 통해 자동으로 최적화된 실행 계획을 생성하며, SQL 쿼리와 DataFrame API를 모두 지원합니다.
SparkSQL
- 구조화된 데이터를 처리하기 위한 모듈
- DataFrame의 핵심 엔진
- SQL처럼 직관적인 데이터 처리 가능
- 자동으로 최적화된 실행 계획 생성
→ RDD의 단점 보완, 분산 SQL 쿼리 엔진 역할 수행
SparkSQL의 역할
- SQL 질의 실행
- 다양한 언어(Java, Scala, Python, R)에서 사용 가능
- 정형화된 포맷(JSON, CSV, Parquet 등) → 임시 테이블 변환
- JDBC/ODBC 커넥터를 통한 외부 연결
- Spark SQL 셸 제공으로 빠른 데이터 탐색 가능
- 최종 실행 최적화 계획 및 JVM 코드 생성

Spark SQL의 내부 작동 구조
- SQL / DataFrame API 입력
- Catalyst Optimizer로 자동 최적화
- 논리 계획 → 물리 계획 변환
- 효율적 실행을 위한 JVM 최적 코드 생성
Catalyst Optimizer
→ Spark SQL의 성능 비결! 내부적으로 복잡한 최적화를 자동 수행

최적화 절차
- SQL Parser & API 해석
- Logical Plan 생성
- Optimized Logical Plan 생성
- Physical Plan 생성
- 최종적으로 RDD 실행 코드로 컴파일
→ 사용자는 코드만 작성하면 Spark SQL이 내부에서 효율적으로 처리!
Catalyst의 최적화 기법
1. Predicate Pushdown
조건 필터를 가능한 한 일찍 적용
# 사용자 코드
df = spark.read.parquet("users.parquet")
result = df.filter(col("age") > 30).select("name", "age")
# Catalyst 최적화: Parquet 파일 읽을 때부터 필터 적용
# → 불필요한 데이터 읽지 않음
2. Column Pruning
필요한 컬럼만 읽기
# 사용자 코드
result = df.select("name").filter(col("age") > 30)
# Catalyst 최적화: name과 age 컬럼만 읽음
# → 메모리 사용량 감소
3. Constant Folding
상수 표현식 미리 계산
# 사용자 코드
df.filter(col("price") > 100 * 1.2)
# Catalyst 최적화
df.filter(col("price") > 120) # 미리 계산
4. Join Reordering
조인 순서 최적화
# 작은 테이블을 먼저 조인하도록 순서 변경
# 중간 결과 크기 최소화
Spark SQL과 DataFrame API의 관계
- 서로 별개가 아님 (같은 Catalyst 엔진 사용)
- 내부적으로 통합된 구조
- DataFrame으로 작성된 명령어도 Spark SQL 엔진을 통해 최적화 및 실행
# 두 방식은 동일한 실행 계획 생성
df.filter(col("age") > 30).select("name")
spark.sql("SELECT name FROM df WHERE age > 30")
Dataset API
- Spark 2.0부터 통합 API로 등장 (DataFrame + Dataset)
- Dataset =
Typed API+Untyped APIDataFrame = Dataset[Row]와 동일
- Java, Scala: Typed API 사용 가능
- Python, R: Type Safety 보장 못해서 DataFrame만 사용 가능

DataFrame vs Dataset
주된 차이: 오류 탐지 시점
- DataFrame: 문법은 Compile-time / 분석은 Runtime
- Dataset: 문법 + 분석 모두 Compile-time
- SQL: 모든 오류 Runtime

오류 탐지 시점 비교
// DataFrame (Untyped)
df.select("name").show() // 컬럼명 오타 → Runtime 오류
// Dataset (Typed - Scala/Java만)
case class Person(name: String, age: Int)
ds.map(_.name).show() // 컬럼명 오타 → Compile 오류
Spark SQL 실행 계획 확인
df = spark.read.parquet("data.parquet")
result = df.filter(col("age") > 30).select("name")
# 논리 계획 확인
result.explain(True)
출력:
== Parsed Logical Plan ==
...
== Analyzed Logical Plan ==
...
== Optimized Logical Plan ==
...
== Physical Plan ==
...
Spark SQL 사용 예시
1. 테이블 등록 및 쿼리
# DataFrame을 임시 테이블로 등록
df.createOrReplaceTempView("people")
# SQL 쿼리 실행
result = spark.sql("""
SELECT name, age
FROM people
WHERE age > 30
ORDER BY age DESC
""")
result.show()
2. 복잡한 집계
spark.sql("""
SELECT
department,
AVG(salary) as avg_salary,
COUNT(*) as employee_count
FROM employees
GROUP BY department
HAVING COUNT(*) > 10
ORDER BY avg_salary DESC
""").show()
3. 조인
spark.sql("""
SELECT
e.name,
e.salary,
d.department_name
FROM employees e
JOIN departments d ON e.dept_id = d.id
WHERE e.salary > 50000
""").show()
Tungsten 실행 엔진
Spark SQL은 Catalyst 외에도 Tungsten 엔진을 사용:
- 메모리 관리 최적화: JVM 객체 오버헤드 제거
- 캐시 친화적 연산: CPU 캐시 활용 극대화
- 코드 생성: 런타임에 최적화된 바이트코드 생성
Catalyst vs RDD Optimizer
| 특징 | Catalyst (DataFrame/SQL) | RDD |
|---|---|---|
| 최적화 시점 | 쿼리 계획 단계 | 사용자가 직접 |
| 최적화 범위 | 전체 쿼리 | 개별 연산 |
| 타입 정보 활용 | O (스키마 기반) | X |
| 코드 생성 | 자동 | 없음 |
| 성능 | 높음 (자동 최적화) | 사용자 의존 |
→ 쿼리를 어떻게 더 똑똑하게 실행할지 (Catalyst)
→ 어떻게 자원을 효율적이고 나눠서 실행할지 (Transformation & Action)
최적화 권장사항
1. SQL vs DataFrame API
# 복잡한 쿼리는 SQL이 가독성 좋음
spark.sql("""
SELECT ... FROM ... WHERE ... GROUP BY ... HAVING ...
""")
# 간단한 변환은 DataFrame API
df.select("name").filter(col("age") > 30)
2. Broadcast Join
작은 테이블 조인 시 Broadcast 힌트 사용
from pyspark.sql.functions import broadcast
large_df.join(broadcast(small_df), "key")
3. 파티션 관리
# 조인 전 적절한 파티션 수 설정
spark.conf.set("spark.sql.shuffle.partitions", "200")