Apache Spark 4.1의 Real-Time Mode와 선언적 데이터 파이프라인

Apache Spark 4.1의 Real-Time Mode, Spark Declarative Pipelines, Spark Connect 구조와 도입 시 고려할 제약을 정리한다.

2026-08-14 · 최초 발행 2026-04-17

스트리밍 지연과 파이프라인 관리 방식을 함께 바꾼 Spark 4.1

Apache Spark 4.1은 Structured Streaming에 공식 Real-Time Mode(RTM)를 도입하고, Spark Declarative Pipelines(SDP)를 추가한 릴리스다. RTM은 마이크로배치 경계를 없앤 연속 실행으로 p99 단일 자리수 밀리초 지연을 목표로 한다. SDP는 데이터셋과 쿼리 선언을 바탕으로 의존 관계와 실행 그래프를 구성한다.

여기에 Spark Connect의 gRPC 기반 클라이언트-서버 모델이 결합된다. Python과 SQL의 선언적 파이프라인, 다국어 클라이언트 통합, 저지연 실행 엔진이 각각 따로 기능하는 것이 아니라 하나의 데이터 처리 환경에서 맞물린다.

RTM은 Structured Streaming에서 long-running task를 유지하는 실행 모드다. SDP는 streaming table, materialized view, flow의 선언을 해석해 실행 순서를 구성하는 프레임워크다. Spark Connect는 DataFrame 연산을 미해결 논리 계획으로 직렬화해 gRPC로 전달하는 분리형 아키텍처다.

마이크로배치 대신 지속 실행하는 Real-Time Mode

기존 Structured Streaming의 마이크로배치는 수백 밀리초~수 초 단위의 트리거마다 쿼리를 스케줄링한다. 배치가 바뀔 때마다 태스크를 다시 시작하는 비용도 발생한다.

RTM은 태스크를 수명이 긴 long-running stage로 유지한다. 레코드가 들어오면 배치 경계를 기다리지 않고 처리하므로, stateless projection·filter·UDF 파이프라인에서 p99 단일 자리수 밀리초 지연을 달성한다. 저지연 피처 엔지니어링에서는 기존 Flink 대비 동등 이상의 성능이 보고되었다.

Spark 4.1 기준으로 RTM은 Scala를 우선 지원하며 PySpark 확장은 차기 릴리스 예정이다. 소스는 Apache Kafka이고, 싱크는 Kafka Sink와 ForeachSink를 지원한다. 적용 대상은 projection·filter·UDF로 이뤄진 stateless 단일 스테이지 쿼리다. spark.sql.streaming.realTimeMode.enabled 설정으로 전환할 수 있어 코드 재작성은 필요하지 않다.

Kafka SourceContinuous Reader(long-running task)Stateless Transform(Projection / Filter / UDF)Continuous WriterKafka Sink / ForeachSinkCheckpoint Store

소스와 싱크 태스크는 실행 중 계속 살아 있으며, 체크포인트 저장소에는 오프셋이 주기적으로 기록된다. 이 구조가 연속 처리와 장애 복구를 함께 지탱한다.

선언한 데이터 객체에서 실행 그래프를 만드는 SDP

SDP에서 pipeline은 개발과 실행의 기본 단위다. 하나의 pipeline 안에는 여러 flow, streaming table, materialized view가 포함될 수 있다. streaming table은 스트리밍 소스에서 지속적으로 증가하는 테이블이고, materialized view는 선언한 쿼리를 바탕으로 자동 갱신되는 구체화 뷰다. flow는 테이블이나 뷰로 레코드를 공급하는 단방향 데이터 이동 단위다.

파이프라인은 Python(.py)과 SQL(.sql)로 작성할 수 있다. 사용자는 만들고자 하는 테이블·뷰·쿼리를 선언하고, SDP는 정의된 객체 사이의 의존 관계를 분석한다. 실행 순서와 병렬화 전략은 런타임이 결정하며, 오케스트레이션·컴퓨트 관리·오류 처리도 이 계층에서 담당한다.

Execution LayerDeclarative LayerSQL / Python 파이프라인 정의SDP Planner: 의존성 분석DAG 생성: streaming table /MV / flowCatalyst Optimizer실행 엔진: 배치 또는 RTM 전환

Declarative Layer는 사용자의 정의를 해석하고 의존성을 파악한다. Execution Layer는 생성된 DAG를 최적화한 뒤 배치 또는 RTM 실행 엔진을 선택한다. 연속 실행과 증분 처리 최적화는 차기 버전에서 추가될 예정이다.

Spark Connect가 분리하는 클라이언트와 실행 환경

기존 JVM 중심 모델에서는 클라이언트 애플리케이션과 드라이버 JVM이 같은 프로세스에서 실행된다. IDE나 애플리케이션 서버와 통합하기 어렵고, 클라이언트는 특정 Spark 버전에 종속된다. 다국어 지원도 PySpark 외에는 제한적이었다.

Spark Connect는 DataFrame API 호출을 미해결 논리 계획으로 인코딩한 뒤 protocol buffers로 직렬화한다. 이 계획은 gRPC 기반 양방향 스트리밍 채널을 통해 Spark Connect Server로 전달된다. 클라이언트는 앱 서버·노트북·IDE에 임베드할 수 있는 thin client가 되고, 분석·최적화·실행은 서버 측 Spark Driver에서 수행된다.

Server Side (Spark Driver)Client Side (thin)Client: Python / Scala / Go/ Rust / SwiftDataFrame APIUnresolved Logical Plan(protobuf)gRPC ChannelSpark Connect ServerAnalyzer / Optimizer /ExecutorCluster Resources

이 모델은 클라이언트의 경량화와 서버 측 실행을 분리한다. Spark 4.1에서는 안정성과 다국어 클라이언트 지원이 한층 강화되었으며, Go·Rust·Swift 커뮤니티 클라이언트도 확장되고 있다.

Spark 4.1 전후의 실행 모델 차이

항목 Spark 4.0 이하 Spark 4.1
스트리밍 최저 지연 수백 ms~수 초 단일 자리수 ms (RTM)
파이프라인 정의 명령형 SQL/DataFrame 선언형 SDP(Python·SQL)
다국어 클라이언트 PySpark/Scala 중심 Go·Rust·Swift까지 확장
실행 그래프 관리 사용자 오케스트레이션 SDP 자동 의존성 분석
RTM 소스/싱크 해당 없음 Kafka / Kafka·Foreach

저지연 피처 처리와 Lakehouse 파이프라인

금융·광고 도메인에서는 RTM을 실시간 피처 스토어 업데이트에 적용할 수 있다. 기존 마이크로배치가 500ms 이상 지연을 만들던 구간을 10ms 이하로 줄이는 사례가 있으며, SQL 친화성과 Spark 생태계 통합성은 Flink와 비교할 때의 강점이다.

Lakehouse에서는 Bronze→Silver→Gold 레이어를 SDP로 선언할 수 있다. streaming table과 materialized view를 함께 구성하고, 팀은 데이터 품질 룰과 기대값을 정의한다. 운영자가 수작업으로 다루던 DAG 재시작과 재처리 로직은 런타임 기본 기능으로 제공된다.

전환 전에 확인할 제약과 운영 범위

RTM은 stateless 쿼리를 우선 지원한다. stateful 조인과 집계는 차기 릴리스 대상이다. Scala 우선 지원이라는 언어 제약도 있어, Python 확장 로드맵을 확인해야 한다.

기존 embedded 드라이버 방식의 애플리케이션은 Spark Connect로 전환할 때 마이그레이션 비용이 발생한다. SDP 역시 명령형 패턴에서 선언적 사고방식으로 옮겨가는 학습 과정이 필요하다.

운영 관점에서는 메트릭과 로깅이 기존 Structured Streaming UI와 동일하므로, 사용 중인 모니터링 대시보드를 재활용할 수 있다. RTM의 지원 범위와 Spark Connect 전환 비용, SDP의 운영 모델을 함께 검토해야 한다.

Spark 4.1은 저지연 스트리밍, 선언적 파이프라인, 다국어 클라이언트라는 축을 같은 방향으로 밀어붙인다. RTM은 마이크로배치 경계를 없애 SQL·DataFrame 생태계에 Flink 수준의 지연을 가져오고, SDP는 데이터 엔지니어링을 만들 대상의 선언으로 압축한다. Spark Connect까지 포함하면 기존 Spark 자산을 유지하면서 실시간·선언적·언어 중립 아키텍처로 점진 전환할 수 있다.

Sources

Apache SparkStructured Streaming데이터 파이프라인Spark Connect실시간 처리