LLM WikiAccess-protected knowledge portal
← 스터디 홈
101편 · 약 12분

Apache Spark 4.2: Auto CDC·Metric Views·지리 타입으로 데이터 엔지니어링의 수작업을 엔진으로 옮긴 방법

Spark 4.2가 해결하려는 문제

Apache Spark 4.0(2025년 초)이 Python UDF Arrow 전환과 VARIANT 타입으로 기초를 다졌다면, 4.2는 그 위에 세 가지 수작업 영역을 엔진 수준으로 흡수한다.

CDC 조립: Debezium 이벤트를 Kafka에서 읽고, 삭제·수정·삽입을 구분해 타깃 테이블에 반영하는 파이프라인을 팀마다 직접 구현한다. 실수가 발생하기 쉽고, 포맷(Debezium, Maxwell, MySQL binlog)마다 처리 로직이 달라진다.

메트릭 정의 분산: 같은 "월간 활성 사용자(MAU)" 지표가 BI 대시보드, 데이터 사이언스 노트북, ML 피처 파이프라인에서 서로 다른 쿼리로 구현돼 있다. 한 곳에서 정의를 바꾸면 나머지가 어긋난다.

지리 데이터 처리: PostGIS나 BigQuery GIS처럼 Spark SQL에 네이티브 지리 타입과 공간 함수가 없어, 외부 라이브러리(Sedona 등)를 별도로 설치하고 관리해야 했다.

Spark 4.2.0은 2026년 7월에 릴리스됐고, Databricks Runtime 19가 이 버전을 기반으로 한다.

Auto CDC: 변경 데이터 캡처를 엔진이 처리한다

CHANGES 절

Spark 4.2는 SQL에 CHANGES 절을 추가한다. Data Source V2 커넥터가 지원하는 경우, 테이블의 행 수준 변경 이력을 배치 및 스트리밍 두 방식으로 읽을 수 있다.

-- 배치: 특정 버전 구간의 변경 이력 읽기
SELECT * FROM my_table CHANGES FROM VERSION 10 TO VERSION 20;

-- 타임스탬프 기준으로도 읽기 가능
SELECT * FROM orders CHANGES FROM TIMESTAMP '2026-07-01' TO TIMESTAMP '2026-07-31';

반환되는 행에는 _change_type 컬럼이 포함되며 값은 insert, update_preimage, update_postimage, delete 중 하나다. 커넥터(Delta Lake, Iceberg 등)가 버전·타임스탬프를 자체 도메인으로 해석하고, Spark는 통합 문법과 후처리를 담당한다.

스트리밍에서도 같은 문법으로 변경 이력을 지속적으로 소비할 수 있다.

-- 스트리밍 CDC 싱크 생성
CREATE STREAMING TABLE cdc_target AS
SELECT *
FROM STREAM source_table CHANGES FROM VERSION 0;

Auto CDC (SCD Type 1 자동 적용)

Auto CDC는 CHANGES API 위에 SCD(Slowly Changing Dimension) Type 1 자동 upsert를 제공하는 파이프라인 레벨 기능이다. 변경 이벤트를 어떻게 타깃 테이블에 반영할지 선언하면 Spark가 복잡한 merge 로직을 생성한다.

from pyspark.sql import SparkSession

spark = SparkSession.builder.getOrCreate()

# Auto CDC — Python 선언형 API
(
  spark.readStream
    .option("readChangeFeed", "true")
    .table("source_orders")
  .apply_changes(
    target="target_orders",
    keys=["order_id"],
    sequence_by="updated_at",         # 이 컬럼 기준으로 최신 이벤트 선택
    apply_as_deletes="operation = 'D'",
    except_column_list=["_metadata"]
  )
)

apply_changes가 내부적으로 MERGE INTO 문을 생성해 실행한다. 이전에는 팀마다 직접 작성하던 MERGE 로직을 표준 API로 대체한다.

Metric Views: 비즈니스 지표를 SQL의 일급 객체로

Metric Views는 Spark SQL에 네이티브 시맨틱 레이어를 추가한다. 팀이 측도(measure)차원(dimension)을 한 번 정의하면, 다양한 도구와 쿼리에서 동일한 정의로 접근한다.

-- Metric View 생성
CREATE METRIC VIEW monthly_active_users
  MEASURES (
    COUNT(DISTINCT user_id) AS mau,
    SUM(session_count)      AS total_sessions
  )
  DIMENSIONS (event_date, country, platform)
  FROM user_events
  WHERE event_type = 'session_start';

-- 어디서든 동일하게 사용
SELECT country, mau
FROM monthly_active_users
WHERE event_date BETWEEN '2026-07-01' AND '2026-07-31';

Spark 엔진은 Metric View를 파싱할 때 집계 의미론을 보존하므로, 동일한 지표를 BI 대시보드, ML 피처 파이프라인, LLM 데이터 인터페이스에서 쓸 때 값이 일치한다. 이전에는 이 일관성을 dbt나 Looker 같은 외부 도구에 의존하거나 팀 규약으로 관리했다.

지리 타입: 공간 분석이 Spark SQL 표준으로 편입

Spark 4.2는 GEOMETRYGEOGRAPHY 타입을 Spark SQL에 네이티브로 추가한다. OGC(Open Geospatial Consortium) 표준 공간 함수를 별도 라이브러리 없이 쓸 수 있다.

-- 지리 타입 사용 예
CREATE TABLE store_locations (
  store_id  STRING,
  location  GEOMETRY
);

-- WKT(Well-Known Text)에서 변환
INSERT INTO store_locations VALUES
  ('S001', ST_GeomFromText('POINT(127.0276 37.4979)'));  -- 서울 역삼

-- 두 지점 간 거리 계산 (미터)
SELECT store_id, ST_Distance(location, ST_GeomFromText('POINT(126.9780 37.5665)')) AS dist_m
FROM store_locations
ORDER BY dist_m;

-- 폴리곤 내 점 검색
SELECT store_id
FROM store_locations
WHERE ST_Within(location, ST_GeomFromText('POLYGON((...))'));

이 기능으로 Apache Sedona 같은 서드파티 지리 라이브러리를 별도로 관리할 필요가 줄어든다. 다만 Sedona의 고급 공간 인덱스(HNSW 기반 공간 조인 최적화 등)는 여전히 외부 라이브러리가 필요하다.

Arrow-first PySpark: 기존 코드가 빨라진다

Spark 4.2는 Arrow 기반 Python UDF와 Arrow-based IPC를 기본값으로 전환한다. 이전 버전에서 옵션으로 켜야 했던 설정이 자동 활성화된다.

from pyspark.sql.functions import udf
from pyspark.sql.types import DoubleType

# 추가 설정 없이 Arrow-first UDF 실행
@udf(DoubleType())
def normalize_score(raw: float) -> float:
    return (raw - 0.5) / 0.5

# 기존 코드 그대로 — 내부에서 Arrow IPC 사용
df.withColumn("score_normalized", normalize_score(df.raw_score))

Arrow IPC를 통해 Python ↔ JVM 직렬화 오버헤드가 크게 줄어든다. 코드 변경 없이 기존 PySpark 파이프라인의 Python UDF 구간 성능이 향상된다.

Real-Time Mode for PySpark: PySpark 스트리밍에서 Python UDF의 지연을 줄이는 새 실행 모드다. 배치 지향적인 Arrow 마이크로배치 대신, 레코드 단위로 더 빠르게 Python 함수를 호출하는 경로를 제공한다.

Spark Connect 확장: 에이전트와 외부 서비스에서 Spark 호출

Spark Connect는 클라이언트와 Spark 서버를 gRPC + Arrow로 분리하는 원격 연결 레이어다. 4.2에서 Spark Connect의 성숙도가 높아져 다음 시나리오가 안정적으로 지원된다.

  • Jupyter / VSCode 노트북에서 원격 Spark 클러스터로 연결
  • 에이전트(Python 프로세스)가 런타임에 Spark 쿼리를 시작
  • CI/CD 파이프라인에서 Spark 잡을 원격으로 제출
from pyspark.sql import SparkSession

# Spark Connect로 원격 Spark 서버에 연결
spark = SparkSession.builder \
    .remote("sc://spark-server.internal:15002") \
    .getOrCreate()

# 이후 일반 PySpark 코드와 동일
df = spark.table("orders").filter("amount > 1000")

클라이언트 사이드에는 Python 패키지만 필요하고, Spark 런타임은 서버에서 실행된다. 조직은 Spark 실행 환경을 중앙에서 관리하면서 사용자는 경량 클라이언트로 접근한다.

기타 주요 변경 사항

VARIANT 타입 개선: is_valid_variant(expr) 함수 추가. 입력이 유효한 VARIANT 값인지 boolean으로 반환한다. JSON 파싱 실패 행을 필터링할 때 유용하다.

Upsert 문법 추가: INSERT INTO ... REPLACE ON/USING 문법이 SQL 표준으로 추가됐다. 키 기준으로 기존 행을 교체하거나 새 행을 삽입하는 upsert를 Delta Lake 없이도 선언적으로 표현할 수 있다.

INSERT INTO products REPLACE ON (product_id)
SELECT product_id, name, price FROM incoming_products;

Spark Web UI 현대화: Bootstrap 5 전환, 다크 모드 지원, SQL 실행 계획 시각화 개선, 쿼리 타임라인 업데이트, 서버사이드 페이지네이션.

Kubernetes 지원 강화: 이기종 executor 관리(GPU executor와 CPU executor를 동일 클러스터에서 혼용), 안정적인 Resource Manager API, 컨트롤 플레인 오버헤드 감소.

Spark 4.2 주요 기능 아키텍처

데이터 수집
Auto CDC
CHANGES 절 + apply_changes()
SCD Type 1 자동 upsert
INSERT REPLACE ON
선언형 upsert
키 기반 병합
저장·처리 레이어
VARIANT + is_valid_variant
반정형 데이터 처리
잘못된 JSON 필터링
GEOMETRY / GEOGRAPHY
ST_* 네이티브 공간 함수
OGC 표준 준수
분석·시맨틱 레이어
Metric Views
MEASURES + DIMENSIONS 정의
일관된 비즈니스 지표
집계 의미론 보존
BI / ML / LLM 동일 값
중복 정의 제거
실행·접근 레이어
Arrow-first PySpark
UDF / IPC 기본값
코드 변경 없이 가속
Spark Connect (gRPC)
원격 접근·에이전트 호출
경량 클라이언트 지원
공통 기반
Spark Web UI 현대화
Bootstrap 5 · 다크 모드
Kubernetes 이기종 Executor
GPU + CPU 혼합 클러스터
Data Source V2 트랜잭션
멀티테이블 원자적 쓰기
Spark 4.2 핵심 기능 구조

업그레이드 고려 사항

기존 PySpark 파이프라인: Arrow-first UDF가 기본 활성화된다. 기존 UDF가 Arrow와 호환되지 않는 타입(예: Python datetime 객체를 반환하면서 타임존 처리가 다른 경우)을 쓰는지 확인한다. 비호환 시 spark.sql.execution.arrow.pyspark.enabled=false로 임시 비활성화 후 단계적으로 마이그레이션한다.

Spark Connect 네트워크 정책: Spark Connect를 활성화하면 포트 15002가 열린다. 클러스터 방화벽 규칙에 이 포트를 명시적으로 허용하거나 차단해야 한다.

Metric Views와 기존 뷰의 공존: Metric View는 일반 뷰와 이름 공간을 공유한다. 기존 뷰 이름과 충돌하지 않도록 네이밍 컨벤션을 정한다.

CHANGES 절 커넥터 지원: CHANGES 절은 Data Source V2 커넥터가 SupportsRowLevelOperations를 구현해야 동작한다. Delta Lake는 지원하고, Iceberg는 버전에 따라 확인이 필요하다. 레거시 Parquet 직접 쓰기 테이블은 지원하지 않는다.

References