구조적 스트리밍 프로그래밍 가이드

구조적 스트리밍 프로그래밍 가이드 (Structured Streaming Programming Guide)

구조적 스트리밍(Structured Streaming)은 Spark SQL 엔진 위에 만들어진 확장 가능하고 장애에 강한(fault-tolerant) 스트림 처리 엔진이에요. 정적 데이터에 대한 배치 계산을 표현하는 것과 같은 방식으로 스트리밍 계산을 표현할 수 있습니다. 이 가이드에서는 프로그래밍 모델과 API를 차근차근 살펴보며, 스트리밍 워드 카운트 예제부터 시작할게요.

출처: 문서

본문

개요 (Overview)

구조적 스트리밍(Structured Streaming)은 Spark SQL 엔진 위에 구축된 확장 가능하고 장애에 강한(fault-tolerant) 스트림 처리 엔진입니다. 정적 데이터에 대한 배치 계산을 표현하는 것과 같은 방식으로 스트리밍 계산을 표현할 수 있어요. Spark SQL 엔진이 이를 점진적이고 연속적으로 실행하고, 스트리밍 데이터가 계속 도착함에 따라 최종 결과를 갱신해 줍니다. Scala, Java, Python 또는 R에서 Dataset/DataFrame API를 사용해 스트리밍 집계, 이벤트-시간 윈도우, 스트림-배치 조인 등을 표현할 수 있어요. 계산은 동일한 최적화된 Spark SQL 엔진 위에서 실행됩니다. 마지막으로 이 시스템은 체크포인트와 Write-Ahead Logs를 통해 종단 간 정확히 한 번(exactly-once) 장애 허용 보장을 제공합니다. 요컨대, 구조적 스트리밍은 사용자가 스트리밍에 대해 고민할 필요 없이 빠르고, 확장 가능하며, 장애에 강하고, 종단 간 정확히 한 번 처리되는 스트림 처리를 제공합니다.

내부적으로 구조적 스트리밍 쿼리는 기본적으로 마이크로-배치 처리(micro-batch processing) 엔진을 사용해 처리되는데, 이 엔진은 데이터 스트림을 일련의 작은 배치 작업으로 처리하여 종단 간 지연 시간을 최대 100밀리초까지 낮추면서 정확히 한 번 장애 허용을 보장해요. 그러나 Spark 2.3부터는 **연속 처리(Continuous Processing)**라는 새로운 저지연 처리 모드가 도입되었는데, 최소 한 번(at-least-once) 보장으로 종단 간 지연 시간을 최대 1밀리초까지 낮출 수 있어요. 쿼리의 Dataset/DataFrame 연산을 바꾸지 않고도 애플리케이션 요구 사항에 따라 모드를 선택할 수 있습니다.

이 가이드에서는 프로그래밍 모델과 API를 함께 살펴볼 거예요. 대부분의 개념은 기본 마이크로-배치 처리 모델을 기준으로 설명하고, 이후에 연속 처리 모델에 대해 논의할게요. 먼저 구조적 스트리밍 쿼리의 간단한 예제인 스트리밍 워드 카운트부터 시작해 봅시다.

더 알아보기 (Learn more)