공부를 하다/Databricks

Day 08. DataFrameReader API + Views

Banaaan 2026. 8. 24. 12:49

Delta 테이블만 읽다 보면 CSV, JSON, Parquet 같은 원천 파일을 읽을 때 막힌다. 오늘은 다양한 파일 포맷을 읽는 DataFrameReader API와, 쿼리를 저장해두는 View 3종류를 정리한다.


DataFrameReader API

pandas의 read_csv()처럼 파일을 읽어 DataFrame으로 만드는 도구다.

spark.read.format("형식").option("옵션명", "옵션값").load("경로")

pandas와 차이점:

  pandas read_csv Spark DataFrameReader
처리 위치 단일 머신 메모리 클러스터 분산 처리
파일 크기 한계 수백 MB 수백 GB
실행 시점 즉시 Lazy (Action 전까지 대기)
폴더 통째로 읽기 불가 가능
# pandas는 파일 하나씩
pd.read_csv("2024-01.csv")

# Spark는 폴더 통째로 — 파일 100개도 한 번에
spark.read.format("csv").option("header", "true").load("/data/2024/")

CSV 읽기

df = spark.read.format("csv") \
    .option("header", "true") \
    .option("inferSchema", "true") \
    .load("/Workspace/Users/practice-user/data/sample.csv")
옵션 역할 안 쓰면?
header 첫 줄을 컬럼명으로 사용 첫 줄이 데이터로 들어옴
inferSchema 컬럼 타입 자동 감지 전부 String으로 읽힘

inferSchema는 편리하지만 파일 전체를 스캔해서 느리다. 실무에서는 스키마를 직접 지정하는 게 낫다.


스키마 직접 지정

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

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

df = spark.read.format("csv") \
    .option("header", "true") \
    .schema(schema) \
    .load("/Workspace/Users/practice-user/data/sample.csv")

inferSchema 대신 .schema(schema)를 쓰면 스캔 없이 바로 읽는다.


JSON 읽기

df = spark.read.format("json") \
    .load("/Workspace/Users/practice-user/data/sample.json")

JSON은 키 이름이 파일 안에 있어서 header 옵션이 필요 없다. inferSchema도 기본으로 켜져 있다.


Parquet 읽기

df = spark.read.format("parquet") \
    .load("/Workspace/Users/practice-user/data/sample.parquet")

Parquet은 스키마 정보가 파일 내부에 저장되어 있어서 옵션이 거의 없다.


축약 문법

spark.read.csv("경로")
spark.read.json("경로")
spark.read.parquet("경로")

.format() 없이 메서드명으로 바로 지정할 수 있다. 시험에 두 형태 모두 나온다.


Views — 쿼리를 저장해두는 방법

View는 테이블처럼 보이지만 실제 데이터가 저장된 건 아니다. "이 SQL 실행하면 이 결과 보여줘"를 이름 붙여 저장한 것이다.

View 3종류 비교 — 시험 핵심

종류 유지 기간 접근 범위
Temp View 세션 종료 시 삭제 현재 세션만
Global Temp View 클러스터 종료 시 삭제 같은 클러스터의 모든 세션
Permanent View 영구 저장 누구나 (메타스토어에 저장)

Temp View

df.createOrReplaceTempView("temp_orders")

spark.sql("SELECT * FROM temp_orders WHERE amount > 100")

createOrReplace — 같은 이름이 있으면 덮어쓴다. 노트북 재실행 시 오류 방지.

세션이 끊기면 사라진다. 나 혼자 쓰는 임시 뷰.


Global Temp View

df.createOrReplaceGlobalTempView("global_orders")

# 조회 시 반드시 global_temp. 접두사 필요
spark.sql("SELECT * FROM global_temp.global_orders")

global_temp.가 없으면 오류난다. 이게 Temp View와 가장 큰 차이.

같은 클러스터에 다른 노트북이 있어도 접근 가능하다. 클러스터 꺼지면 사라진다.


Permanent View

spark.sql("""
    CREATE OR REPLACE VIEW databricks_practice.default.vw_high_orders
    AS SELECT * FROM databricks_practice.default.orders WHERE amount > 100
""")

메타스토어에 저장되어 클러스터를 껐다 켜도 남아있다. 팀원도 접근 가능하다.

이름 앞에 vw_를 붙이는 건 테이블과 구분하기 위한 관습이다.


한 줄 요약

Temp View        → 내 세션만, 껐다 켜면 없어짐
Global Temp View → 클러스터 살아있는 동안, global_temp. 필수
Permanent View   → 영구 저장, 메타스토어에 등록

JOIN

df_orders.join(df_customers, on="customer_id", how="inner")
how 결과
inner 양쪽 다 있는 것만
left 왼쪽 기준, 오른쪽 없으면 null
right 오른쪽 기준, 왼쪽 없으면 null
outer 양쪽 모두, 없으면 null

집계 (Aggregation)

from pyspark.sql.functions import count, sum, avg, max, min

df.groupBy("category") \
  .agg(
      count("order_id").alias("주문수"),
      sum("amount").alias("총매출"),
      avg("amount").alias("평균금액")
  )

결과 컬럼명이 sum(amount)처럼 지저분하게 나오므로 항상 .alias()로 이름을 붙인다.


Window Functions

groupBy 집계는 행이 줄어들지만, Window는 행 수를 유지하면서 각 행에 계산값을 추가한다.

from pyspark.sql.functions import rank, row_number, lag, lead
from pyspark.sql.window import Window

window_spec = Window.partitionBy("category").orderBy("amount")

df.withColumn("순위", rank().over(window_spec))
함수 특징
rank() 동점이면 같은 순위, 다음 순위 건너뜀 (1,1,3)
row_number() 동점 없이 무조건 1,2,3
lag("컬럼", n) n행 이전 값
lead("컬럼", n) n행 이후 값

✓ Day 1: Spark 기초
✓ Day 2: Delta Lake
✓ Day 3: Auto Loader
✓ Day 4: Structured Streaming
✓ Day 5: Delta Live Tables
✓ Day 6: Workflows
✓ Day 7: Architecture + Notebooks + Git Folders
✓ Day 8: DataFrameReader API + Views