콘텐츠로 이동

구조적 스트리밍 (Structured Streaming)

Spark Structured Streaming 은 정적 배치와 같은 방식으로 스트리밍 데이터를 표현하는 API 입니다. 스트리밍을 배치처럼(DataFrame/Dataset) 작성하면 엔진이 점진·내결함성 있게 실행합니다. 이 장은 개념을 데이터스케쳐스의 실시간 수집 관점에서 정리합니다.

배치처럼 쓰는 스트리밍

Structured Streaming 의 핵심은 "정적 데이터로 쓰던 변환을 그대로 스트리밍에 쓰게 하는 것"입니다.

stream = spark.readStream.format("kafka").option("kafka.bootstrap.servers", "...").load()
result = stream.selectExpr("CAST(value AS STRING) AS msg").filter("msg IS NOT NULL")
query = result.writeStream.outputMode("append").format("console").start()
  • 지연(late) 데이터는 워터마크로 처리합니다.
  • 정확히 한 번(exactly-once) 처리를 위한 폴트-톨러런트 상태 저장을 엔진이 담당합니다.

무한 스트리밍 테이블 개념

  • 입력·출력을 무한히 커지는 테이블로 봅니다.
  • outputMode 로 append / update / complete 를 골라 출력 방식을 정합니다.
  • 트리거 간격으로 마이크로 배치(또는 연속 처리)를 실행합니다.

데이터스케쳐스 실무 관점

  • 이벤트 데이터 수집(D-SKET Events): 현장/행사에서 들어오는 이벤트 로그를 Kafka → Spark Structured Streaming 으로 실시간 집계(방문·참여 지표)하거나, 배치 워크로드와 같은 변환 로직을 재사용할 때 적합합니다.
  • 배치 ETL 과 스트리밍을 같은 DataFrame API 로 통합하면, 로직을 두 번 구현하지 않아도 됩니다.

확인 필요

  • 스파크 버전(4.x)에 따라 가이드가 세분화되어 있으니, 실제 사용 버전의 Structured Streaming 문서를 재확인하세요. (확인 필요)

더 알아보기