LLM WikiAccess-protected knowledge portal

WIKI

Apache Spark 4.1: VARIANT 쉬레딩 GA, Real-Time Mode, SQL 스크립팅으로 스트리밍과 분석을 다시 쓰는 방법

왜 Spark 4.1인가 Apache Spark 4.0이 ANSI 기본값 전환, VARIANT 타입 도입, Java 17 기본 환경 등 4.x 세대의 틀을 세웠다면, 4.1은 그 기반 위에서 세 가지 핵심 기능을 운영 가능한 수준으로 완성 했다. 1. VARIANT 쉬레딩 Shredding GA 반정형 semi structured 데이터를 이진 블롭에 넣던 방식에서 벗어나, 자주 쓰는 필드를 타입드 Parquet 컬럼으로 분리

경로human/study/content/database-frontier/16-spark-4-1-variant-shredding-rtm-sql-scripting.md
카테고리Study
태그#mysql #rtm #scripting #shredding #sql #study #variant

왜 Spark 4.1인가

Apache Spark 4.0이 ANSI 기본값 전환, VARIANT 타입 도입, Java 17 기본 환경 등 4.x 세대의 틀을 세웠다면, 4.1은 그 기반 위에서 세 가지 핵심 기능을 운영 가능한 수준으로 완성했다.

  1. VARIANT 쉬레딩(Shredding) GA: 반정형(semi-structured) 데이터를 이진 블롭에 넣던 방식에서 벗어나, 자주 쓰는 필드를 타입드 Parquet 컬럼으로 분리해 읽기를 최대 30배 빠르게 한다.
  2. Structured Streaming Real-Time Mode(RTM): 마이크로배치 구조를 버리고 연속 데이터 흐름으로 전환해 P99 레이턴시를 40~300ms로 낮춘다. 상태 없는 작업은 한 자릿수 밀리초도 가능하다.
  3. SQL 스크립팅 GA: 오류 핸들러, 다중 변수 선언 등 프로시저 수준 SQL 제어가 기본값으로 활성화된다.

이 세 변화는 각각 다른 고통 지점에 닿는다. 반정형 데이터 분석, 실시간 스트리밍, 복잡한 SQL 로직을 별도 도구에 위탁하던 이유가 하나씩 줄어든다.

릴리스 일정 참고: Spark 4.1.0은 2026년 봄에 공개됐으며, 4.1.2 패치가 2026년 5-6월에 출시됐다. 정확한 4.1.0 릴리스 날짜는 공식 Apache Spark news 페이지에서 확인할 것.


Apache Spark 4.1 — 세 가지 핵심 기능 아키텍처 ① VARIANT 쉬레딩 GA 원본 반정형 데이터 (JSON/이벤트 로그) {"event":"click","user_id":42,"ts":1720000000} Parquet 쉬레딩 metadata (이진) value (전체 fallback) typed_value 컬럼 ↓ event:varchar / user_id:bigint / ts:bigint 쿼리 I/O 절감 읽기 최대 30× 빠름 · 쓰기 20~50% 느림 row group/page 스킵, 컬럼 프루닝 적용 ② Structured Streaming RTM 마이크로배치 (기존) 배치 경계 기다림 RTM (신규) 연속 데이터 흐름 RTM 세 가지 혁신 연속 데이터 흐름 도착 즉시 처리 배치 경계 없음 파티션당 장기 실행 태스크 파이프라인 스케줄링 스테이지 동시 실행 업스트림 완료 대기 없음 즉시 다운스트림 스트리밍 셔플 디스크 없이 태스크 간 메모리 직접 전달 셔플 병목 제거 ③ SQL 스크립팅 GA DECLARE 다중 변수 선언 DECLARE x INT, y INT DEFAULT 0; CONTINUE HANDLER (오류 복구) DECLARE CONTINUE HANDLER FOR SQLEXCEPTION ... 기본값 활성화 4.1부터 SQL 스크립팅 기본 켜짐 Python 없이 복잡한 파이프라인 가능 Arrow-native UDF / UDTF @pandas_udf → @arrow_udf 데코레이터 Pandas 변환 제거 → 직접 PyArrow 처리 Python Data Source 필터 푸시다운 지원 청크 단위 Arrow 결과 스트리밍 Spark Declarative Pipelines (SDP) 데이터셋 + 쿼리 선언 → Spark가 실행 그래프 관리 의존성 순서·병렬성·체크포인트·재시도 자동화 Python/SQL 혼용, Spark Connect 멀티언어 RTM과 네이티브 통합 (초저지연 파이프라인) 기타 성능 개선 zstd 압축 protobuf 플랜 (대형 워크로드) KLL·Theta 스케치 (근사 집계 GA) 재귀 CTE 지원 대형 로컬 릴레이션 안정성 강화 Spark 4.1이 없애는 외부 의존성들 반정형 분석에 별도 JSON 파싱 레이어 · 초저지연 스트리밍에 Flink 별도 운영 · 파이프라인 로직에 Python 오케스트레이터
Apache Spark 4.1 핵심 기능 흐름

1. VARIANT 쉬레딩: 반정형 데이터를 칼럼 스토어처럼 읽기

문제: 이진 블롭에 갇힌 JSON

Spark 4.0이 도입한 VARIANT 타입은 JSON 같은 반정형 데이터를 스키마 정의 없이 컬럼 하나에 담는 방법이다. 편리하지만 기존 방식에는 한계가 있었다. 쿼리할 때 전체 이진 블롭을 파싱해야 하므로, event 하나 필드만 필요해도 레코드 전체를 읽어야 했다. 칼럼형 스토리지의 I/O 최적화 장점이 사라지는 문제다.

VARIANT 쉬레딩의 동작 원리

쉬레딩은 VARIANT 컬럼의 자주 쓰는 필드를 별도 타입드(typed) Parquet 컬럼으로 추출해 저장하는 기법이다. Parquet 파일 내부에서 VARIANT 그룹은 세 부분으로 구성된다.

파트역할
metadata (이진)타입 정보, 딕셔너리
value (이진)전체 VARIANT 직렬화 값 (fallback)
typed_value 그룹쉬레딩된 서브 컬럼들 (예: event=varchar, user_id=bigint)

쿼리 엔진은 typed_value 그룹에서 필요한 서브 컬럼만 골라 읽는다. 전체 이진 블롭을 파싱하지 않아도 된다.

쉬레딩의 성능 수치와 트레이드오프

결론적으로 읽기가 주가 되는 분석 테이블에는 쉬레딩이 명확히 유리하다. 반면 빠른 쓰기가 중요하고 조회가 드문 로그 수집 파이프라인이라면 쉬레딩을 끄거나 대상 컬럼을 줄이는 게 낫다.

운영 시 확인 포인트

-- 쉬레딩 대상 필드를 지정해 테이블 생성
CREATE TABLE events (
  ts TIMESTAMP,
  payload VARIANT
) USING DELTA
TBLPROPERTIES ('delta.feature.variantType-preview' = 'supported');

-- VARIANT 쉬레딩 힌트 (Spark SQL)
SELECT payload:event, payload:user_id
FROM events
WHERE payload:user_id = 42;

쉬레딩된 테이블에서 payload:event처럼 경로 표기로 필드를 참조하면, Spark는 typed_value.event만 읽는 플랜을 생성한다.


2. Structured Streaming Real-Time Mode(RTM): 마이크로배치 장벽 제거

왜 마이크로배치가 한계에 부딪히는가

Spark Structured Streaming의 기존 마이크로배치 모드는 수백 밀리초~수 초 단위로 소규모 배치를 반복 실행한다. 처리량 면에서는 충분히 강하지만, 레이턴시 하한이 배치 간격에 묶인다. 결제 사기 탐지, 운영 알림, 실시간 집계 같이 수십 밀리초 응답이 필요한 시나리오에서는 Flink나 별도 스트리밍 엔진을 추가하거나, Spark를 버려야 했다.

RTM은 Spark Structured Streaming API를 바꾸지 않고 이 레이턴시 하한을 깬다.

RTM의 세 가지 엔진 혁신

연속 데이터 흐름(Continuous Data Flow)

RTM은 파티션당 하나의 장기 실행 태스크를 유지한다. 태스크는 소스(Kafka 등)에서 레코드를 도착하는 대로 즉시 처리한다. 배치 경계가 없으므로 배치 시작까지 기다리는 시간이 제거된다.

파이프라인 스케줄링(Pipeline Scheduling)

마이크로배치에서는 스테이지가 직렬로 실행된다. 스테이지 A가 완전히 끝나야 스테이지 B가 시작한다. RTM에서는 여러 스테이지가 동시에 실행된다. 업스트림 레코드가 스테이지 A를 통과하는 즉시 다운스트림 스테이지 B가 해당 레코드를 처리한다.

스트리밍 셔플(Streaming Shuffle)

기존 셔플은 중간 결과를 디스크에 쓰고 다음 스테이지가 읽는다. 디스크 I/O가 레이턴시 병목이 된다. RTM의 스트리밍 셔플은 태스크 간 데이터를 메모리에서 직접 전달해 디스크 접근을 없앤다.

RTM이 달성하는 레이턴시

워크로드 유형P99 레이턴시
무거운 운영 워크로드40~300ms
상태 없는(stateless) 작업한 자릿수 ms

적합한 사용 사례: 사기 탐지, 실시간 집계, 알림 파이프라인, 이벤트 기반 트리거

적합하지 않은 사례: 대용량 배치, 복잡한 조인이 많은 ETL(마이크로배치가 처리량 면에서 유리)

RTM과 Spark Declarative Pipelines 연계

RTM은 Spark Declarative Pipelines(SDP)와 네이티브로 통합된다. SDP에서 파이프라인을 선언하고 실행 모드를 RTM으로 지정하면, 파이프라인 스케줄링, 체크포인트, 재시도를 Spark가 자동 관리한다. 별도 오케스트레이터 없이 서브초 파이프라인을 운영할 수 있다.


3. SQL 스크립팅 GA: 프로시저 제어를 SQL로

Spark 4.0에서 프리뷰로 들어온 SQL 스크립팅이 4.1에서 GA가 되고 기본값으로 활성화됐다. 두 가지 개선이 특히 실용적이다.

다중 변수 DECLARE

-- 4.0 이전: 변수마다 개별 선언
DECLARE x INT DEFAULT 0;
DECLARE y INT DEFAULT 0;

-- 4.1: 한 줄로 여러 변수
DECLARE x INT, y INT DEFAULT 0;

대형 SQL 스크립트에서 선언부를 줄이고 가독성을 높인다.

CONTINUE HANDLER for 정교한 오류 복구

BEGIN
  DECLARE CONTINUE HANDLER FOR SQLEXCEPTION
  BEGIN
    -- 오류 시 fallback 처리
    INSERT INTO error_log SELECT current_timestamp(), sqlerrm();
  END;

  -- 여기서 오류가 나면 HANDLER로 넘어간 뒤 다음 구문 계속
  INSERT INTO target SELECT * FROM source WHERE ...;
  UPDATE metrics SET processed = processed + 1;
END;

CONTINUE HANDLER는 오류 발생 후 스크립트를 중단하지 않고 핸들러 블록을 실행한 다음, 오류가 난 구문 다음부터 계속 진행한다. Python에서 try/except 처럼 쓸 수 있는 SQL 레벨 오류 처리다.


4. Arrow-native UDF/UDTF: Pandas 변환 없이 PyArrow 직접

기존 Pandas UDF(@pandas_udf)는 내부적으로 Spark 내부 형식 → Pandas DataFrame → 처리 → 다시 Spark 내부 형식 변환 경로를 거쳤다. 4.1에서 도입된 Arrow-native UDF(@arrow_udf)는 이 변환 경로를 줄인다.

from pyspark.sql.functions import arrow_udf
import pyarrow as pa
import pyarrow.compute as pc

@arrow_udf(pa.float64())
def normalize(arr: pa.ChunkedArray) -> pa.ChunkedArray:
    mean = pc.mean(arr).as_py()
    std = pc.stddev(arr).as_py()
    return pc.divide(pc.subtract(arr, mean), std)

PyArrow 배열을 받아 PyArrow 배열로 반환한다. Pandas 변환 없이 Arrow 네이티브 연산을 쓸 수 있으므로 CPU 오버헤드가 줄어든다.

Python Data Source 필터 푸시다운: 커스텀 Python Data Source에서 Spark의 필터 조건을 직접 받아 소스 레벨에서 데이터 양을 줄일 수 있다. API, 파일 시스템, 서드파티 연동에서 네트워크 및 파싱 비용을 줄이는 패턴에 유용하다.


5. 기타 안정성·성능 개선

zstd 압축 protobuf 플랜

대형 워크로드에서 쿼리 플랜이 너무 커지면 드라이버 메모리 압박과 직렬화 지연이 발생한다. 4.1은 플랜을 zstd로 압축해 네트워크 전송과 메모리 사용을 줄였다. 수백 스테이지, 수천 태스크 규모의 파이프라인에서 효과가 두드러진다.

KLL·Theta 스케치 GA

근사 집계(Approximate Aggregation)에서 정확한 COUNT(DISTINCT) 대신 확률적 자료구조로 빠르게 근사값을 구하는 패턴이다.

대규모 분석 쿼리에서 근사 집계로 레이턴시를 낮추고 싶을 때 쓸 수 있다.

재귀 CTE 지원

WITH RECURSIVE employee_hierarchy AS (
  SELECT id, name, manager_id, 0 AS level
  FROM employees
  WHERE manager_id IS NULL
  UNION ALL
  SELECT e.id, e.name, e.manager_id, h.level + 1
  FROM employees e
  JOIN employee_hierarchy h ON e.manager_id = h.id
)
SELECT * FROM employee_hierarchy;

계층 구조 데이터 탐색에 재귀 CTE가 SQL 표준으로 지원된다.


운영자 전환 체크리스트

VARIANT 쉬레딩 도입 시

RTM 도입 시

SQL 스크립팅 GA 시

Arrow-native UDF 전환


Open question

References