본문으로 건너뛰기

Databricks Feature Store의 200ms 실시간 피처

Spark RTM과 Lakebase로 Kafka 이벤트를 200ms p99의 온라인 피처로 전환합니다.

이 요약은 AI가 원문을 분석해 생성했습니다. 정확한 내용은 원문 기준으로 확인하세요.

TL;DR

사기 탐지와 개인화 모델은 장기 사용자 기준선만으로는 최근 행동 변화를 놓치므로 Kafka에서 들어오는 이벤트를 밀리초 단위로 피처에 반영해야 합니다. Databricks Feature Store는 Spark Real-Time Mode가 각 이벤트를 연속 처리하고 RocksDB에서 롤링 윈도 상태를 갱신한 뒤 Lakebase에 최신 값을 쓰도록 구성하며, Model Serving은 추론 시 해당 피처를 자동으로 가져옵니다. RTM은 마이크로배치 대신 동시 실행과 분산된 체크포인트를 사용하고, Lakebase는 소규모 upsert의 WAL 쓰기 증폭을 낮춰 이 경로를 뒷받침합니다. Kafka 이벤트부터 온라인 피처 가용성까지 end-to-end p99 latency는 200ms이며, 장애 시 최대 5분 데이터를 재생하는 방식으로 exactly-once 처리를 유지합니다.

섹션별 상세

01
사기 탐지와 개인화 모델은 수십 일간의 사용자 기준선과 최근 몇 분간의 행동 신호를 함께 봐야 현재 상황에 맞는 판단을 내릴 수 있습니다. 기존 배치 Spark 파이프라인은 주기적으로 집계값을 갱신하므로 피처에 수분에서 수시간의 지연이 생겼고, 밀리초 단위 신호가 필요한 경우 별도 스트리밍 로직과 운영 인프라가 필요했습니다. Databricks Feature Store는 하나의 피처 정의를 오프라인 배치와 온라인 실시간 파이프라인에서 함께 사용하도록 구성해 이 인프라 부담을 줄이고, Kafka 이벤트 도착부터 온라인 피처 저장까지 end-to-end p99 latency 200ms를 목표로 합니다.
02
Kafka 이벤트는 serverless Lakeflow Spark Delta Pipelines의 Spark Real-Time Mode (RTM) 파이프라인으로 들어가고, 여기서 롤링 집계가 계산됩니다. 갱신된 값은 streaming JDBC sink를 통해 Lakebase 온라인 피처 저장소에 기록되며, Model Serving endpoint가 추론 시 최신 피처를 조회해 모델 입력과 결합합니다. 예를 들어 사용자별 최근 10분 거래액은 로컬 RocksDB의 상태값을 읽고 새 거래액을 더한 뒤 Lakebase에 기록하므로, 새 결제 승인 요청에 과거 구매 기준선과 최신 거래 합계가 함께 제공됩니다.
Kafka의 사용자 이벤트가 Spark RTM 파이프라인과 RocksDB 로컬 상태를 거쳐 Lakebase 온라인 Feature Store에 저장되고, Model Serving endpoint가 이를 추론에 사용하는 흐름을 나타낸 아키텍처 다이어그램입니다.
Diagram다이어그램은 이벤트 수집, 로컬 상태 갱신, 실시간 집계, 온라인 저장, 추론 시 조회라는 처리 경로를 왼쪽에서 오른쪽으로 연결합니다. 특히 RocksDB가 집계 상태를 빠르게 유지하고 Lakebase가 저지연 조회용 값을 제공하는 분리가 RTM 기반 실시간 피처 파이프라인의 핵심 구조로 나타납니다.
03
시간 윈도는 갱신 빈도와 최신성 사이의 선택을 결정합니다. Tumbling window는 12:00–12:10처럼 고정된 경계에서 한 번 갱신되고, Sliding window는 5분 간격으로 겹치는 구간을 생성하며, Rolling window는 각 이벤트 시점에서 과거 10분을 다시 계산해 새 이벤트마다 값을 갱신합니다. 따라서 자주 변하지 않는 피처에는 업데이트 횟수와 운영 비용이 적은 Tumbling 또는 Sliding window가 적합하고, 모든 이벤트를 즉시 반영해야 하는 실시간 사기 신호에는 Rolling window가 적합합니다.
고정 경계에서 갱신되는 Tumbling window와 각 이벤트마다 이동하며 갱신되는 Rolling window의 시간 집계 동작을 비교한 다이어그램입니다.
Diagram상단의 Tumbling window는 10분 경계에 도달할 때만 결과를 갱신해 경계 이전의 최신 이벤트가 즉시 반영되지 않는 구조를 나타냅니다. 하단의 Rolling window는 이벤트가 발생할 때마다 직전 10분 구간을 이동시켜 새 집계값을 만들며, 원문이 설명한 실시간 사기 탐지와 최신 사용자 의도 포착에 필요한 즉시성을 시각화합니다.
04
기존 Spark Structured Streaming의 Microbatch Mode (MBM)는 이벤트를 배치로 모은 뒤 각 단계를 순차 처리하고 배치 경계마다 체크포인트를 기록해 상태 집계 지연의 하한을 만듭니다. RTM은 데이터 검증·변환 단계와 엔터티별 집계 단계를 동시에 실행하며, 각 행이 도착하는 즉시 RocksDB 상태를 갱신하고 새 집계값을 downstream으로 내보냅니다. 윈도 만료 시점에는 해당 이벤트의 기여분을 제거한 수정값도 다시 기록하고, 체크포인트 비용은 더 긴 구간의 행들에 분산해 정상 처리 지연을 낮춥니다.
05
RTM 파이프라인은 장애 복구를 위해 exactly-once 처리를 유지하면서 최대 5분의 Kafka 데이터를 재생할 수 있도록 설계됐습니다. serverless Lakeflow Spark Delta Pipelines가 클러스터 프로비저닝과 용량 계획을 맡고, 인프라 업데이트 때 새 클러스터가 준비된 뒤 기존 클러스터를 중단해 피처 최신성의 중단을 최소화합니다. 5분 체크포인트 간격에 맞춘 핸드오프로 유지보수 중 재처리 공백과 서비스 중단을 줄이는 구조입니다.
06
Lakebase는 compute와 distributed storage를 분리해 온라인 피처 조회와 빈번한 소규모 upsert를 처리합니다. 표준 Postgres는 자주 갱신되는 행에서 첫 변경 때 전체 8KB 페이지 이미지를 WAL에 기록해 쓰기 증폭이 커질 수 있지만, Lakebase는 분산 safekeeper 노드의 quorum이 작은 변경 레코드를 확인하도록 해 쓰기 경로의 로그량을 줄입니다. 그 결과 RTM이 이벤트마다 최신 피처를 계속 게시하면서도 추가 지연을 낮출 수 있고, 온라인 Feature Store는 수만 reads per second와 수십 ms 수준의 지연을 처리하도록 확장됩니다.
07
Model Serving은 추론 시 Lakebase에서 필요한 피처를 자동 조회해 별도의 조회 코드나 수동 연결 작업을 없앱니다. MLflow로 모델을 기록할 때 피처 의존성이 함께 저장되고, 요청이 들어오면 Model Serving이 해당 의존성을 바탕으로 최신 집계값을 가져와 추론 요청에 결합합니다. serving 인프라는 CPU endpoint에서 100K+ QPS를 지원하고 트래픽 변화에 맞춰 수평 확장하므로, 실시간 피처 계산과 대규모 추론을 하나의 경로로 연결합니다.
08
스트리밍 피처는 원천 스트림의 보존 기간이 짧아 학습용 과거 데이터 생성이 어려울 수 있습니다. Databricks Feature Store는 Kafka 입력의 오프라인 사본을 저장하고, 과거 시점에 맞는 피처를 다시 계산해 point-in-time join을 수행하며 온라인 스트리밍 피처의 backfill에도 같은 기능을 사용합니다. 또한 피처를 Unity Catalog의 first-class object로 관리하고 접근 제어·lineage·MLflow 모델 의존성을 연결해, 여러 인프라에 흩어진 피처의 재사용과 거버넌스를 단일 플랫폼 안에서 유지합니다.

용어 해설

롤링 윈도 집계(Rolling Window Aggregation)
각 이벤트의 타임스탬프를 기준으로 직전 일정 구간의 합계·평균·개수를 다시 계산하는 방식입니다. 고정된 시계 경계가 아니라 새 이벤트가 들어올 때마다 윈도가 이동하며, 만료된 이벤트의 기여분을 제거해 현재 시점에 가까운 값을 유지합니다. 실시간 사기 탐지나 개인화처럼 최신 이벤트가 즉시 반영되어야 하는 피처에 중요합니다.
마이크로배치 처리(Microbatch Processing)
스트리밍 데이터를 짧은 시간 단위의 배치로 모아 처리하는 실행 방식입니다. 각 배치가 수집되고 여러 처리 단계를 거친 뒤 체크포인트가 기록되므로, 이벤트가 도착한 즉시 집계하지 못하고 배치 주기와 처리 시간만큼 지연이 생깁니다. 원문에서는 Spark Real-Time Mode와 대비되는 기존 방식으로 제시됩니다.
RocksDB 상태 저장소(RocksDB State Store)
상태 기반 스트리밍 연산에서 엔터티별 중간 집계값과 만료 정보를 저장하는 로컬 저장소입니다. 각 executor가 보유한 RocksDB에 이벤트별 상태를 기록하므로 새 행이 들어올 때 값을 즉시 증가시키거나 만료된 기여분을 제거할 수 있습니다. 메모리보다 큰 상태도 처리할 수 있다는 점이 실시간 윈도 집계에 활용됩니다.
쓰기 전 로그(Write-Ahead Log (WAL))
데이터 변경을 먼저 로그에 기록해 장애 복구와 내구성을 보장하는 저장 방식입니다. 표준 Postgres에서는 한 페이지의 첫 변경 때 작은 논리 변경 대신 전체 8KB 페이지 이미지를 기록할 수 있어, 빈번한 소규모 upsert에서 로그량이 커지는 문제가 생깁니다. Lakebase는 분산 저장 구조를 이용해 작은 변경 레코드 중심으로 쓰기 증폭을 줄입니다.
시점 일치 조인(Point-in-Time Join)
학습 시점에 실제로 이용 가능했던 피처 값만 연결하는 데이터 결합 방식입니다. 과거 이벤트와 피처의 유효 시점을 맞춰 미래 정보가 학습 데이터에 섞이는 것을 막고, 스트리밍 피처의 과거 값을 다시 계산할 때도 동일한 시간 기준을 유지합니다. 원문에서는 오프라인 학습 데이터 생성과 온라인 피처 backfill에 사용됩니다.

기술

  • Databricks Feature Store
  • Kafka
  • Spark Real-Time Mode (RTM)
  • Lakeflow Spark Delta Pipelines
  • RocksDB
  • Lakebase
  • streaming JDBC sink
  • Model Serving
  • MLflow
  • Unity Catalog
  • Postgres

활용 사례

  • 실시간 신용카드 사기 탐지
  • 사용자 행동 기반 개인화
  • 클릭스트림과 광고 노출 집계
  • 최근 거래액과 장기 구매 기준선을 결합한 결제 승인
  • 실시간 피처를 사용하는 저지연 모델 추론
AI 분석 전체 내용 보기

AI 요약 · 북마크 · 개인 피드 설정 — 무료

출처 · 인용 안내

원문 발행 2026. 08. 18.수집 2026. 08. 18.출처 타입 RSS

인용 시 "요약 출처: AI Trends (aitrends.kr)"를 표기하고, 사실 확인은 원문 보기 기준으로 진행해 주세요. 자세한 기준은 운영 정책을 참고해 주세요.