RisingWave → Snowflake 실시간 파이프라인 구축하기: 신선한 데이터를 더 빠르게
Snowflake 웨어하우스에 분석 가능한 신선한 데이터를 가져오는 일은 마치 시간과의 경쟁처럼 느껴질 수 있습니다. 신뢰할 수 있는 전통적인 배치 기반 ETL 프로세스는 종종 지연을 초래하며, 이는 인사이트가 항상 실제보다 몇 시간(혹은 그 이상) 뒤처지게 만듭니다. 만약 이 지연 시간을 대폭 줄이고, 이미 변환된 상태로 BI 도구나 쿼리에 바로 사용할 수 있는 데이터가 Snowflake에 도달할 수 있다면 어떨까요?
바로 그것이 RisingWave의 새로운 Snowflake 싱크 커넥터가 가능하게 하는 일입니다.
RisingWave는 이벤트 스트림 데이터를 실시간으로 수집, 처리, 변환할 수 있도록 설계된 통합 실시간 데이터 플랫폼입니다. 데이터가 어디로 가는지뿐 아니라, 어디에서 오는지도 매우 유연하게 다룰 수 있습니다. Apache Kafka, Apache Pulsar와 같은 인기 있는 스트리밍 플랫폼은 물론, PostgreSQL이나 MySQL 같은 트랜잭션 데이터베이스의 CDC(Change Data Capture) 스트림, 또는 애플리케이션과 마이크로서비스 간의 상호작용에서 발생하는 이벤트 스트림까지도 손쉽게 처리할 수 있습니다. RisingWave는 이질적인 실시간 데이터를 SQL 기반으로 복잡한 연산(조인, 집계, 윈도잉 등)을 수행할 수 있는 강력한 엔진이라고 할 수 있습니다.
이제 Snowflake 싱크 커넥터를 통해 RisingWave에서 Snowflake 테이블로 직접 연결되는 지속적인 데이터 파이프라인을 구축할 수 있습니다.
그렇다면 이것이 Snowflake에서의 데이터 활용 방식에 어떤 변화를 가져올까요?
데이터 신선도를 극대화하여 빠르게 실행할 수 있습니다.
주기적인 배치 작업을 기다릴 필요 없이, RisingWave는 처리된 데이터를 지속적으로 스트리밍합니다. Snowflake에 도달하는 데이터는 최신 이벤트를 반영하며, 대시보드나 리포트가 항상 최신 상태를 유지합니다.Snowflake 내 데이터 전처리 부담이 줄어듭니다.
여러 스트림 간 조인, 실시간 집계, 복잡한 필터링 및 변환 작업이 모두 RisingWave 내부에서 처리되기 때문에, Snowflake에 도달하는 데이터는 이미 정제된 상태입니다. 이는 Snowflake 내의 복잡한 변환 논리를 단순화하고 컴퓨팅 비용도 절감할 수 있게 해줍니다.
작동 방식은 어떻게 될까요?
전체 프로세스는 매우 매끄럽고 자동화되어 있습니다. RisingWave는 설정된 소스에서 데이터를 수집하고, 사용자가 정의한 실시간 변환을 수행한 뒤, 처리된 데이터를 JSON 형식으로 사용자가 관리하는 S3 버킷에 저장합니다. 이후 Snowflake의 Snowpipe 서비스가 이 데이터를 해당 테이블로 자동으로 로드합니다. 이 방식은 Snowpipe의 S3 연동 자동화 기능과 완벽히 호환되도록 설계되어 있습니다.
Snowpipe를 쓰는 상황에서 왜 RisingWave가 필요한가요?
좋은 질문입니다. 데이터가 결국 S3에 저장되고 Snowpipe가 Snowflake로 로드하는 일을 한다면, 굳이 RisingWave를 추가해야 할까요?
하지만 결정적인 차이는 데이터가 S3에 도달하기 이전 단계에서 발생합니다:
실시간 복잡 변환 처리
RisingWave는 스트리밍 중인 데이터에 대해 SQL 기반의 조인, 집계, 필터링, 데이터 보강을 실시간으로 수행할 수 있습니다. 반면 Snowpipe는 단순히 데이터를 로드할 뿐이며, 변환은 보통 Snowflake 내에서 후처리로 진행되므로 지연이 발생하고 별도 컴퓨팅 리소스를 요구하게 됩니다.상태 기반 처리 및 머티리얼라이즈드 뷰 지원
RisingWave는 증분 방식으로 저지연 업데이트가 가능한 머티리얼라이즈드 뷰를 유지할 수 있어, 실시간 세션화, 리더보드와 같은 복잡한 분석 뷰를 미리 계산할 수 있습니다. 이렇게 사전 처리된 정제 데이터를 Snowflake로 전송하면, 원시 이벤트가 아닌 실질적인 인사이트가 전달됩니다.다양한 소스와 직접 통합 및 초기 필터링 처리
RisingWave는 Kafka, Pulsar, 데이터베이스 CDC 스트림과 같은 다양한 소스에 직접 연결할 수 있으며, 수집 프로토콜 처리, 역직렬화, 초기 필터링 및 스키마 검증 등을 사전 처리합니다. 반면 Snowpipe만 사용하려면, 이러한 작업을 별도로 처리할 시스템이나 커스텀 코드가 필요합니다.Snowflake 내 처리 단순화 및 지연 감소
Snowflake로 유입되는 데이터가 이미 분석 가능한 상태로 변환된 덕분에, 별도 변환 로직 없이 즉시 쿼리 가능하며, 전체 파이프라인의 응답 속도와 리소스 효율성이 개선됩니다.
시작해 보기: 간단한 설정 예시
설정을 시작하려면 S3 버킷 정보와 Snowpipe가 해당 버킷을 감시하도록 설정되어 있어야 합니다. (관련 가이드는 Snowflake 공식 문서 참고)
RisingWave에서 싱크를 생성하는 방식은 다음과 같이 간단합니다:
CREATE SINK snowflake_sink
FROM ss_mv -- 머티리얼라이즈드 뷰 또는 소스
WITH (
connector = 'snowflake',
type = 'append-only', -- 또는 upsert 지원! 아래 참고
s3.bucket_name = 'your-s3-bucket-name',
s3.credentials.access = 'your-aws-access-key',
s3.credentials.secret = 'your-aws-secret-key',
s3.region_name = 'your-s3-bucket-region',
s3.path = 'path/to/data/', -- 선택 사항
force_append_only = 'true' -- append-only 모드 사용 시
);
이제 RisingWave가 S3 버킷에 데이터를 스트리밍하고, Snowpipe가 이를 자동으로 Snowflake로 로드합니다.
업데이트와 삭제(Upsert)는 어떻게 처리되나요?
Snowflake는 Snowpipe를 통한 직접적인 Upsert를 지원하지 않지만, RisingWave는 이를 우회하는 똑똑한 방법을 제공합니다.
AS CHANGELOG 옵션으로 머티리얼라이즈드 뷰를 정의하면, 삽입뿐 아니라 업데이트/삭제까지 추적 가능합니다. 이를 위해 __op(삽입, 삭제, 업데이트 유형) 및 __row_id(정렬 기준)와 같은 특수 컬럼이 함께 출력됩니다.
CREATE SINK snowflake_sink as WITH sub AS changelog FROM user_behaviors
SELECT
user_id,
target_id,
event_timestamp AT TIME ZONE 'America/Indiana/Indianapolis' as event_timestamp,
changelog_op AS __op, -- 변경 유형
_changelog_row_id::bigint AS __row_id -- 변경 순서를 위한 ID
FROM
sub WITH (
connector = 'snowflake',
type = 'append-only',
s3.bucket_name = 'EXAMPLE_S3_BUCKET',
s3.credentials.access = 'EXAMPLE_AWS_ACCESS',
s3.credentials.secret = 'EXAMPLE_AWS_SECRET',
s3.region_name = 'EXAMPLE_REGION',
s3.path = 'EXAMPLE_S3_PATH',
);
이러한 변경 로그는 S3에 저장되며, 이후 Snowflake의 Dynamic Table 기능을 활용하여 자동으로 병합 처리할 수 있습니다. 예:
-- Snowflake 내에서
CREATE OR REPLACE DYNAMIC TABLE current_user_behaviors
TARGET_LAG = '1 minute'
WAREHOUSE = your_snowflake_warehouse
AS
SELECT *
FROM (
SELECT *,
ROW_NUMBER() OVER (
PARTITION BY primary_key_column
ORDER BY __row_id DESC
) as rn
FROM your_staging_table_fed_by_snowpipe
)
WHERE rn = 1 AND (__op = 1 OR __op = 3); -- 1 = 삽입, 3 = 업데이트 이후
이 Dynamic Table은 최신 상태의 행만을 자동으로 추출하여 항상 일관된 뷰를 제공합니다.
사용 가능성
Snowflake 싱크 커넥터는 RisingWave Premium Edition 기능 중 하나이며, RisingWave Cloud에서는 별도 비용 없이 제공됩니다. 자체 호스팅 시에는 4코어 이하에서는 무료로 사용 가능하며, 그 이상은 라이선스 키가 필요합니다.
데이터에서 인사이트까지, 지금 가속해보세요
오래된 데이터를 기다리는 시대는 끝났습니다. RisingWave의 Snowflake 싱크 커넥터를 통해 신선하고 사전 처리된 데이터를 기반으로 진정한 실시간 분석 환경을 구축해보세요.
RisingWave 시작하기:
전문가 상담:
복잡한 사용 사례나 데모가 필요하다면 문의하기를 통해 도움을 받아보세요.커뮤니티 참여하기:
Slack 커뮤니티에서 질문하고, 다른 개발자와 교류해보세요.
이제 여러분의 실시간 데이터를 통해 Snowflake에서 어떤 새로운 가능성을 열어갈지 기대됩니다!
