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
'공부를 하다 > Databricks' 카테고리의 다른 글
| Day 10. Delta Lake 심화 — 최적화와 데이터 적재 (0) | 2026.08.25 |
|---|---|
| Day 09. Spark UDF + Higher Order Functions (0) | 2026.08.24 |
| Day 07. Databricks 아키텍처 + Notebooks + Git Folders (0) | 2026.08.24 |
| Day 06. Workflows — 파이프라인 자동화의 진화 (0) | 2026.08.01 |
| Day 05. Delta Live Tables — 파이프라인을 선언만 하면 알아서 실행되는 구조 (0) | 2026.08.01 |