실시간 교통 예측 시스템을 직접 구축했습니다 — 전체 데모 코드는 GitHub에서 확인하세요

이 글에서는 RisingWave의 Streaming SQL 기능을 활용해 실시간 교통 예측 시스템을 구축하는 방법을 소개합니다. 이 시스템은 시뮬레이션된 교통 데이터를 수집하고, 실시간으로 피처 엔지니어링을 수행하며, 사전 학습된 LSTM 모델을 통해 교통 흐름을 예측합니다. 모든 결과는 대시보드에서 시각화됩니다. 이 접근 방식은 스트리밍 SQL이 실시간 머신러닝 파이프라인 구축을 얼마나 쉽게 만들 수 있는지를 보여줍니다.
실시간 교통 예측 시스템 아키텍처

시스템의 복잡성을 이해하기 위해 일반적인 실시간 교통 예측 시스템의 구조를 간단히 살펴보겠습니다. 보통 다음과 같은 주요 단계로 구성됩니다:
데이터 수집:
스마트 카메라와 지상 센서에서 발생하는 원시 데이터가 지속적으로 스트리밍됩니다. 이 데이터는 차량 수, 속도, 차량 종류 및 타임스탬프를 포함합니다.스트림 처리 및 피처 엔지니어링:
고속의 원시 데이터를 머신러닝 모델이 이해할 수 있는 의미 있는 신호(피처)로 변환합니다. 예를 들어, 1분 단위의 평균 속도나 특정 구간을 통과하는 차량 수를 계산합니다.ML 모델 추론:
전처리된 피처가 사전 학습된 ML 모델(LSTM 등)에 입력되어 미래 교통량을 예측합니다.예측 제공 및 액션:
예측 결과는 대시보드 업데이트, 신호등 타이밍 조정, 내비게이션 앱에 전달되는 등 다운스트림 애플리케이션으로 전달됩니다.데이터 저장 및 시각화:
피처와 예측 값은 저장되어 모니터링, 추가 분석, 그리고 대시보드에 과거 데이터를 제공하는 데 활용됩니다.
'실시간'이라는 특성은 이 모든 과정을 훨씬 더 까다롭게 만듭니다. 데이터가 단순히 방대하기만 한 것이 아니라 매우 빠르게 유입되기 때문입니다.
"피처 신선도(feature freshness)"는 절대적으로 중요합니다. 데이터가 단 몇 분만 오래되어도 예측 성능은 현저히 떨어집니다. 또한 높은 처리량의 구성 요소를 원활하게 통합하고 데이터 일관성을 유지하는 것은 상당히 복잡한 엔지니어링 도전입니다.
RisingWave로 실시간 교통 예측 데모 살펴보기
이벤트 스트림 데이터 플랫폼이 실시간 머신러닝을 어떻게 혁신하는지 알아보기 위해 RisingWave로 만든 교통 예측 데모를 살펴보겠습니다. 이 데모는 GitHub에서 확인할 수 있으며, 실시간 데이터 수집부터 예측 생성 및 대시보드 시각화까지 전체 파이프라인을 시뮬레이션합니다.


목표: 분 단위 교통량 예측
이 데모의 목표는 여러 도로 구간에 대해 다음 1분 동안의 교통량(차량 수)을 예측하는 것입니다. 이러한 예측은 동적 교통 관리와 스마트 내비게이션을 위한 핵심 요소입니다.
아키텍처: RisingWave를 중심으로 한 스트리밍 데이터 플로우

전체 시스템은 지속적인 데이터 흐름을 중심으로 설계되어 있으며, RisingWave는 다음과 같은 중요한 역할을 수행합니다:
데이터 수집:
Python 스크립트가 시뮬레이션된 스마트 카메라(차량 속도, 차량 종류)와 지상 센서(차량 수) 데이터를 생성하여 Apache Kafka 토픽으로 지속적으로 전송합니다.RisingWave – 실시간 피처 엔지니어링:
RisingWave는 Kafka 토픽에 직접 연결되어 원시 이벤트 스트림을 실시간으로 처리합니다. 주요 작업은 피처 엔지니어링입니다.각 도로별로 분 단위 차량 수(
vehicle_count) 계산각 도로별로 분 단위 평균 속도(
avg_speed_kph) 계산
이 작업은 Materialized View(물리화 뷰)를 통해 수행됩니다. 물리화 뷰는 새로운 데이터가 도착할 때마다 자동으로 증분 업데이트됩니다. 예를 들어 최근 5분 데이터를 제공하는 뷰는 다음과 같이 정의할 수 있습니다:
-- RisingWave에서의 피처 엔지니어링 예시 CREATE MATERIALIZED VIEW features_last_5min AS SELECT road_id, window_start, vehicle_count, avg_speed_kph FROM features -- 'features'는 원시 카운트와 속도를 결합한 또 다른 물리화 뷰입니다 WHERE window_start <= NOW() - INTERVAL '1 minute' AND window_start >= NOW() - INTERVAL '6 minutes';이
features_last_5min뷰는 사실상 온라인 피처 스토어로 작동합니다.ML 모델 예측:
사전 학습된 LSTM 모델이 포함된 Python 스크립트가 RisingWave의features_last_5min뷰를 쿼리합니다. 이를 통해 모델은 항상 최신 피처 데이터로 예측을 수행합니다. 이 스크립트는 데이터 정규화 등의 사전 처리도 담당합니다.RisingWave – 예측 결과 저장 및 대시보드 데이터 제공:
Python 스크립트에서 생성된 예측 결과는 또 다른 Kafka 토픽으로 전송됩니다. RisingWave는 이 예측 데이터를 다시 수집해 입력 피처와 함께 테이블에 저장합니다. 또한 RisingWave 내의 또 다른 물리화 뷰는 최근 3시간 동안의 실제 및 예측 데이터를 준비하여 라이브 대시보드에 효율적으로 제공합니다.시각화:
웹 기반 대시보드는 현재 교통 상황, 과거 트렌드, 그리고 예측된 교통 흐름을 시각화합니다. 이 모든 데이터는 RisingWave에서 실시간으로 제공됩니다.
전체 구현 보기:
이 개요는 RisingWave가 시스템 중심에서 어떤 역할을 수행하는지 보여줍니다. 전체 SQL 스크립트, Python 코드 및 데모 실행을 위한 단계별 가이드는 GitHub 저장소에서 확인할 수 있습니다.
왜 스트리밍 ML에 RisingWave가 최적인가?
이 교통 예측 데모는 실시간 ML 애플리케이션을 구축할 때 RisingWave가 탁월한 선택인 이유를 명확히 보여줍니다:
SQL 기반의 간편함: 복잡한 스트리밍 로직과 피처 엔지니어링(예: 분 단위 집계)을 익숙한 SQL만으로 정의할 수 있어 개발 시간과 복잡도가 대폭 감소합니다.
항상 최신 상태의 피처 제공: Materialized View는 피처(
features_last_5min)를 자동으로 증분 업데이트합니다. 별도의 배치 작업 없이도 ML 모델은 항상 최신의 저지연 데이터에 접근할 수 있습니다.통합된 단일 플랫폼: 데이터 수집, 변환, 피처 제공, 예측 결과 저장, 대시보드 제공까지 모든 작업을 RisingWave 하나로 처리할 수 있습니다.
실시간 성능 최적화: 높은 처리량과 낮은 지연 시간에 최적화되어 있어 라이브 데이터 속도를 따라잡을 수 있으며, 즉각적인 인사이트를 제공합니다.
즉, RisingWave는 복잡한 실시간 데이터 파이프라인을 훨씬 더 쉽고 효율적으로 구축하고 운영할 수 있게 도와줍니다.
시작할 준비가 되셨나요?
전체 데모 코드는 GitHub 저장소에서 확인할 수 있습니다.
