DataFrame API¶
Spark 의 DataFrame 은 이름 있는 컬럼으로 구성된 분산 컬렉션으로, 관계형 테이블(또는 R/Python 의 데이터프레임)과 개념적으로 같지만 내부적으로 강력한 최적화 엔진(Catalyst)을 사용합니다. 데이터스케쳐스 실무 관점에서 정리합니다.
DataFrame 이란¶
- DataFrame =
Dataset[Row]— named columns 를 가진 분산 데이터 구조. - RDD 수준의 저수준 맵/필터 대신 스키마(구조) 정보를 최적화에 활용합니다.
- SQL 쿼리와 같은 실행 엔진(Catalyst)을 공유하므로, DataFrame API 와 SQL 을 섞어 써도 동일하게 최적화됩니다.
만들기와 읽기¶
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("ds").getOrCreate()
# JSON/Parquet 등 구조화 소스에서
df = spark.read.parquet("s3://bucket/events.parquet")
# 기존 RDD 나 로컬 데이터에서
df2 = spark.createDataFrame([("a", 1), ("b", 2)], ["name", "count"])
- Parquet·ORC·Avro·Hive·JDBC 등 다양한 소스를 동일한 방식으로 읽습니다.
주요 변환 (Transformation)¶
result = (df
.filter(df.status == "ok")
.select("user_id", "amount")
.groupBy("user_id")
.sum("amount")
.orderBy("user_id"))
- 변환은 lazy(게으른) — 액션(show/count/collect/write)을 호출하기 전까지는 실행되지 않습니다. 이 덕분에 Catalyst 가 전체 실행 계획을 최적화합니다.
- SQL 로도 같은 결과를 낼 수 있어, 코드와 SQL 을 상황에 맞게 오갑니다.
RDD vs DataFrame vs SQL¶
- RDD: 저수준, 타입 최적화는 없지만 유연한 함수형 조작.
- DataFrame: 스키마 기반, Catalyst 최적화 — 대부분의 구조화 작업에 권장.
- SQL: 익숙한 문법, DataFrame 과 같은 엔진.
데이터스케쳐스 실무 관점¶
- 대용량 이벤트 집계: 수집된 이벤트를 DataFrame 으로 읽어 Parquet 로 저장하고, 일 단위 집계를 SQL/DataFrame 으로 수행합니다. lazy 평가로 계획을 최적화해 셔플(셔플)을 줄입니다.
- 스키마 잡기: 무스키마 데이터는
spark.read.json처럼 스키마를 추론하되, 비용이 크면 명시 스키마를 지정합니다.
확인 필요¶
- 세부 API(옵티미제이션 내역, 셔플 파티션 수 등)는 Spark 버전별로 다르므로 사용 버전 공식 문서를 재확인하세요. (확인 필요)