전체 그래프
Spark

DataFrame

data-engineeringsparkdataframestructured-data

상위: Spark

요약

DataFrame은 스키마를 가진 분산 데이터 컬렉션으로, RDD의 한계를 보완합니다. 행과 열로 구성된 표 형식이며, Catalyst Optimizer를 통해 자동 최적화됩니다. Pandas DataFrame과 유사하지만 분산 처리가 가능합니다.

RDD의 문제점

RDD API의 한계

  • RDD API는 연산과 표현식을 Spark가 검사하거나 최적화할 수 없음
    • 어떤 연산이 일어나는지 Spark가 알 수 없음
    • PySpark에서는 Iterator[T] 형태로 받아들이기 때문에 데이터 타입을 인식하지 못함
  • 데이터 압축 기술도 적용 불가
    • 타입 정보 부족으로 컬럼 단위 접근이나 바이너리 최적화가 어려움
  • 결과적으로 효율적인 질의 계획 불가
  • 스키마(schema) 정보가 없음
    • 데이터 값 표현은 가능하지만, 메타데이터나 스키마에 대한 명시적 표현 방법이 없음

DataFrame API

  • 스키마(schema)를 가진 분산 데이터 컬렉션
    • **행(row)**과 **열(column)**로 구성된 표 형식
    • 각 열은 명확한 데이터 타입과 메타데이터를 가짐
  • RDD의 한계를 보완하는 구조화된 데이터 모델 (Spark SQL)

DataFrame 특징

  • Pandas DataFrame과 유사함
  • 컬럼 이름 + 스키마를 가진 인메모리 테이블처럼 동작
  • 사람이 보는 형태는 일반적인 표(table)

DataFrame

Data Type

  • 기본 타입: Byte, Short, Integer, Long, Float, Double, String, Boolean, Decimal
  • 정형 타입: Binary, Timestamp, Date, Array, Map, Struct, StructField
  • 📄 Spark 공식 문서

Schema

외부 데이터 소스를 읽을 때 명시적 스키마를 정의하면 여러 이점:

  • 타입 추론 필요 없음 → 속도 증가
  • 잘못된 타입 조기 발견
  • 파일 전체를 읽지 않아도 됨

스키마 정의 예시

from pyspark.sql.types import StructType, StructField, StringType, IntegerType

schema = StructType([
    StructField("name", StringType(), True),
    StructField("age", IntegerType(), True),
    StructField("city", StringType(), True)
])

df = spark.read.csv("data.csv", schema=schema)

DataFrame 구성 요소

DataFrame = RDD + Schema + DSL

  • Schema → Named Columns with Types: 열마다 이름과 타입을 가짐
  • DSL (Domain-Specific Language): .filter(), .select() 등 다양한 고수준 API 제공

DSL vs SQL

# DSL 방식
sciDocs = data.filter(col("label") == 1)
-- SQL 방식
spark.sql("SELECT * FROM data WHERE label = 1")

RDD vs DataFrame

구분RDDDataFrame
데이터 표현 방식값만 표현 가능, 스키마 표현 불가명확한 스키마(컬럼명, 타입)를 가진 구조적 데이터
최적화 및 성능최적화 어려움, 직접 연산 필요Catalyst Optimizer로 자동 최적화 및 빠른 처리
사용 편의성낮음 (저수준 API)높음 (고수준 API, SQL 활용 가능)

→ DataFrame은 메타 정보를 활용하여 효율적이고 최적화된 분석 가능

RDD vs DataFrame 선택

RDD를 언제 쓰고 DataFrame을 언제 쓰는지는 RDD의 사용 케이스 참고.

DataFrame을 사용하는 경우

  1. 도메인 기반 고수준 API가 필요할 때
  2. filter, map, agg, SQL 등 고수준 표현 필요할 때
  3. 타입 안정성과 최적화된 실행 계획 필요할 때
  4. Catalyst, Tungsten(고성능 실행 엔진) 기반의 효율적인 코드 제너레이션 필요할 때
  5. 일관성 있는 Spark API 사용을 원할 때

DataFrame 생성

1. 데이터로부터 생성

from pyspark.sql import Row

# RDD로부터 생성
parts = spark.sparkContext.parallelize([("Mine", "28"), ("Filip", "29")])
people = parts.map(lambda p: Row(name=p[0], age=int(p[1])))
df = spark.createDataFrame(people)

2. 파일로부터 생성

# CSV
df = spark.read.csv("data.csv", header=True, inferSchema=True)

# JSON
df = spark.read.json("data.json")

# Parquet
df = spark.read.parquet("data.parquet")

# 텍스트
df = spark.read.text("data.txt")

3. 스키마와 함께 생성

schema = StructType([
    StructField("name", StringType(), True),
    StructField("age", IntegerType(), True)
])

df = spark.createDataFrame(data, schema)

DataFrame 구조 변환

# DataFrame  RDD
rdd1 = df.rdd

# DataFrame  JSON 문자열 ( 항목 확인)
df.toJSON().first()

# DataFrame  Pandas
pandas_df = df.toPandas()
print(pandas_df)

→ toPandas는 메모리 기반이므로 큰 데이터셋에는 주의

DataFrame 기본 조작

데이터 확인

df.dtypes            # 컬럼과 데이터 타입
df.show()            # 데이터 출력
df.head()            #   조회
df.first()           #  번째 Row 객체

df.take(n)           # n개의  반환
df.schema            # 스키마 확인
df.describe().show() # 요약 통계
df.columns           # 컬럼 이름

head() return single row object / head(n) return list

중복 확인 및 제거

df.count()                     # 전체  
df.distinct().count()         # 고유  
df.printSchema()              # 스키마 출력
df.explain()                  # 실행 계획 (논리, 물리)
df.dropDuplicates()          # 중복 제거

View 등록 및 SQL 실행

# View로 등록
df.createTempView("viewName")  # 현재 세션
df.createGlobalTempView("viewName")  # 모든 세션
df.createOrReplaceTempView("viewName")  # 덮어쓰기 가능

# SQL 실행
spark.sql("SELECT * FROM viewName").show()
spark.sql("SELECT * FROM global_temp.viewName").show()

DataFrame 최적화 팁

1. 스키마 명시

스키마 명시(inferSchema 회피)는 Data IO 참고.

2. 파티션 관리

# 적절한 파티션  유지
df = df.repartition(200)  # 너무 작은 파일 방지
df = df.coalesce(10)      # 파티션  줄이기

3. 필요한 컬럼만 선택

# Bad: 모든 컬럼 로드
df = spark.read.parquet("data.parquet")
result = df.filter(col("age") > 30).select("name")

# Good: 필요한 컬럼만 로드
result = spark.read.parquet("data.parquet") \
              .select("name", "age") \
              .filter(col("age") > 30)

4. 캐싱

# 여러  사용하는 DataFrame은 캐싱
df.cache()
df.count()  #  실행: 캐싱
df.show()   #  번째 실행: 캐시에서 읽음