공부를 하다/Databricks

Day 04. Structured Streaming — 데이터가 흘러오는 걸 실시간으로 처리하기

Banaaan 2026. 7. 28. 12:49

Auto Loader는 새 파일이 오면 처리하는 방식이었다. Structured Streaming은 데이터 스트림 자체를 실시간으로 처리한다.


Auto Loader vs Structured Streaming

Auto Loader:          파일 단위로 감지
                      file1.json → file2.json → file3.json

Structured Streaming: 데이터가 흘러오는 스트림을 계속 처리
                      row1 → row2 → row3 → row4 → ...

실무에서 어떤 걸 쓰느냐는 데이터 특성에 따라 달라진다.

방식 언제
Structured Streaming (Kafka) 클릭, 주문, 결제 등 내부 서비스 이벤트 실시간 처리
Auto Loader 협력사 파일, 정산 데이터 등 외부에서 주기적으로 들어오는 파일

Kafka란?

여러 시스템에서 데이터가 동시에 쏟아질 때 중간에서 받아서 전달해주는 메시지 큐다.

[Producer] 쇼핑몰 클릭 / 주문 이벤트 / 결제 완료
      ↓
[Kafka Topic] 메시지 쌓임
      ↓
[Consumer] Databricks / 알림 서비스 / 데이터 웨어하우스

Databricks DE 입장에서는 "Kafka Topic에서 데이터 읽어서 Delta Table에 쓰는 파이프라인"을 만드는 것이 주 역할이다.


Delta Table을 스트리밍 소스로 쓸 수 있는 이유

Delta Table은 _delta_log에 모든 변경 이력을 기록한다. Structured Streaming이 이 로그를 보고 새 데이터를 감지한다.

orders_source에 새 데이터 추가
        ↓
_delta_log에 변경 기록
        ↓
Streaming이 로그 보고 새 데이터 감지
        ↓
orders_result에 append

실습 코드

소스 테이블 생성

data = [
    (1, "order_001", 15000),
    (2, "order_002", 32000),
    (3, "order_003", 8000)
]

df = spark.createDataFrame(data, ["id", "order_id", "amount"])
df.write.format("delta").saveAsTable("orders_source")

스트리밍으로 읽고 쓰기

df_stream = spark.readStream \
    .format("delta") \
    .table("orders_source")

df_stream.writeStream \
    .format("delta") \
    .option("checkpointLocation", "/Workspace/Users/practice-user/stream_checkpoint") \
    .trigger(availableNow=True) \
    .toTable("orders_result")

변환 추가

from pyspark.sql.functions import col

df_stream = spark.readStream \
    .format("delta") \
    .table("orders_source")

df_transformed = df_stream.withColumn("amount_with_tax", col("amount") * 1.1)

df_transformed.writeStream \
    .format("delta") \
    .option("checkpointLocation", "/Workspace/Users/practice-user/stream_checkpoint2") \
    .trigger(availableNow=True) \
    .toTable("orders_with_tax")

스트리밍으로 읽으면서 변환까지 한 번에 처리된다.


Auto Loader와 구조 비교

  Auto Loader Structured Streaming
소스 파일 Delta Table, Kafka 등
포맷 cloudFiles delta, kafka
readStream → writeStream 동일 동일
checkpointLocation 필요 필요

소스만 다르고 구조는 동일하다. 하나를 이해하면 나머지도 자연스럽게 연결된다.


오늘로 Associate 기초 파트의 핵심인 Delta Lake, Auto Loader, Structured Streaming을 모두 실습했다. 다음은 Delta Live Tables(DLT)다.