📊 데이터공학

RisingWave

클라우드 네이티브 스트리밍 데이터베이스

상세 설명

RisingWave는 2021년 설립된 RisingWave Labs가 Rust로 개발한 클라우드 네이티브 스트리밍 데이터베이스입니다. PostgreSQL 와이어 프로토콜을 지원해 기존 PostgreSQL 클라이언트와 도구로 스트리밍 데이터를 SQL로 쿼리합니다.

핵심 개념은 Materialized View입니다. CREATE MATERIALIZED VIEW로 스트리밍 데이터에 대한 집계, 조인, 윈도우 쿼리를 정의하면, 새 데이터가 도착할 때마다 결과가 자동으로 증분 갱신됩니다. 전체 재계산 없이 변경분만 처리해 실시간 분석이 가능합니다.

Kafka, Redpanda, Kinesis, S3, PostgreSQL CDC 등 다양한 소스에서 데이터를 수집하고, PostgreSQL, Kafka, Elasticsearch 등으로 결과를 내보냅니다. Source(입력)와 Sink(출력)를 SQL DDL로 정의해 ETL 파이프라인을 코드 없이 구성합니다.

클라우드 네이티브 아키텍처로 컴퓨팅과 스토리지가 분리되어 있고, S3 호환 오브젝트 스토리지를 영구 저장소로 사용합니다. Kubernetes 환경에서 수평 확장이 용이하며, Apache Flink나 ksqlDB 대안으로 주목받고 있습니다.

코드 예제

-- Kafka Source 생성 (이벤트 스트림 입력)
CREATE SOURCE user_events (
    user_id VARCHAR,
    event_type VARCHAR,
    page_url VARCHAR,
    event_time TIMESTAMPTZ,
    properties JSONB
) WITH (
    connector = 'kafka',
    topic = 'user-events',
    properties.bootstrap.server = 'redpanda:9092',
    scan.startup.mode = 'latest'
) FORMAT PLAIN ENCODE JSON;

-- 기본 테이블 (디멘션 데이터)
CREATE TABLE users (
    user_id VARCHAR PRIMARY KEY,
    name VARCHAR,
    email VARCHAR,
    created_at TIMESTAMPTZ
);

-- Materialized View: 실시간 페이지뷰 집계
CREATE MATERIALIZED VIEW page_views_per_minute AS
SELECT
    DATE_TRUNC('minute', event_time) AS minute,
    page_url,
    COUNT(*) AS view_count,
    COUNT(DISTINCT user_id) AS unique_users
FROM user_events
WHERE event_type = 'page_view'
GROUP BY 1, 2;

-- Materialized View: 사용자 세션 분석 (윈도우 함수)
CREATE MATERIALIZED VIEW user_sessions AS
SELECT
    user_id,
    event_time,
    event_type,
    LAG(event_time) OVER (PARTITION BY user_id ORDER BY event_time) AS prev_event_time,
    CASE
        WHEN event_time - LAG(event_time) OVER (PARTITION BY user_id ORDER BY event_time) > INTERVAL '30 minutes'
        THEN 1 ELSE 0
    END AS new_session_flag
FROM user_events;

-- Materialized View: 스트림-테이블 조인
CREATE MATERIALIZED VIEW enriched_events AS
SELECT
    e.user_id,
    e.event_type,
    e.event_time,
    u.name AS user_name,
    u.email
FROM user_events e
LEFT JOIN users u ON e.user_id = u.user_id;

-- Sink: 결과를 Kafka로 출력
CREATE SINK alerts_sink FROM (
    SELECT * FROM page_views_per_minute
    WHERE view_count > 1000  -- 비정상 트래픽 감지
)
WITH (
    connector = 'kafka',
    topic = 'traffic-alerts',
    properties.bootstrap.server = 'redpanda:9092'
) FORMAT PLAIN ENCODE JSON;

-- PostgreSQL 호환 쿼리 (MV 결과 조회)
SELECT * FROM page_views_per_minute
WHERE minute >= NOW() - INTERVAL '5 minutes'
ORDER BY view_count DESC
LIMIT 10;

-- psql로 연결
-- psql -h localhost -p 4566 -d dev -U root

-- Python에서 연결 (psycopg2)
-- import psycopg2
-- conn = psycopg2.connect(
--     host="localhost", port=4566,
--     database="dev", user="root"
-- )
-- cursor = conn.cursor()
-- cursor.execute("SELECT * FROM page_views_per_minute LIMIT 10")

실무에서 이렇게 말해요

시니어: "Flink 대신 RisingWave 쓰면 SQL로만 스트리밍 파이프라인 구성 가능해요. Java 코드 안 짜도 돼요."

주니어: "PostgreSQL 클라이언트로 연결되나요?"

시니어: "네, psql, DBeaver 다 돼요. 기존 BI 도구도 연동 가능하고요."

면접관: "실시간 데이터 처리 경험을 말씀해주세요."

지원자: "RisingWave로 실시간 대시보드 백엔드를 구축했습니다. Kafka 이벤트를 Materialized View로 집계해 1초 이내 지연으로 메트릭을 제공하고, PostgreSQL 프로토콜 덕분에 기존 Metabase를 그대로 연결해 개발 비용을 줄였습니다."

리뷰어: "MV에서 조인할 때 Primary Key가 있어야 효율적으로 증분 갱신돼요. Source에 PRIMARY KEY 설정하세요."

개발자: "user_id를 Primary Key로 설정하면 되겠네요."

주의사항

더 배우기