Snowpipe Streaming
Snowpipe Streaming
Snowpipe Streaming은 Snowflake의 최신 고성능 아키텍처 위에 구축된 실시간 수집(ingestion) 서비스예요. 애플리케이션이 장치, 애플리케이션, 서비스에서 Snowflake 테이블 또는 Snowflake 관리 Apache Iceberg 테이블로 행을 직접 스트리밍할 수 있게 해줍니다. 이 직접 경로 덕분에 워크로드가 필요로 하지 않는 스테이징 파일, 중간 객체 스토리지, 메시지 버스, 커넥터 서비스를 제거할 수 있어요.
출처: Snowflake 문서
본문
Snowpipe Streaming은 두 가지 수집 모드를 지원합니다. 두 모드 모두에서 채널(channel)은 pipe를 통해 대상 테이블로 행을 운반하는 논리적 경로예요:
- Elastic Channels — 대부분의 새 애플리케이션에 권장되는 시작점입니다. 프로듀서가 채널을 만들거나 조정하지 않고 직접 작성합니다. Snowflake가 채널을 관리하고 트래픽 변화에 따라 수집을 확장해요. 승인(acknowledgement)은 Snowflake가 append를 내구성 있게 버퍼링했음을 확인해 주므로 프로듀서는 보유한 복사본을 해제할 수 있고, 이어서 테이블 처리와 쿼리 가시성이 따라옵니다. Elastic Channels는 순서 보장 없이 at-least-once 전달을 제공합니다. 프로듀서는 자신의 전달 요구사항에 따라 미승인 이벤트를 보관합니다.
- Named Channels — offset token을 사용해 각 채널 내에서 순서 있는 exactly-once 수집을 제공합니다. Kafka 파티션이나 CDC(변경 데이터 캡처)처럼 엄격한 순서 의미론이 필요한 소스에서 읽을 때 Named Channels를 사용하세요.
Snowpipe Streaming이 제공하는 것:
- 테이블당 최대 20 GB/s 처리량
- 최저 5초의 수집-조회 가능 지연 시간
- Elastic Channels를 통한 애플리케이션 관리 채널 없이 직접 수집
- Named Channels와 offset token을 통한 순서 있는 exactly-once 수집
- Snowflake 관리 Apache Iceberg 테이블로의 스트리밍
관찰되는 결과는 워크로드 형태와 구성(행 크기, 테이블 폭(컬럼 수), SDK 버퍼링 또는 REST 요청 일괄 처리, 동시성, 테이블 유형, 변환, 클러스터링)에 따라 달라져요.
왜 Snowpipe Streaming을 사용하나요
- 더 단순한 직접 수집: Elastic Channels 덕분에 프로듀서가 장치와 서비스에서 Snowflake로 행을 직접 스트리밍할 수 있어, 채널을 만들거나 프로듀서 간에 수집을 조정할 필요 없이 파이프라인 홉이 줄어듭니다. Snowflake는 프로듀서와 트래픽이 변함에 따라 수집 경로를 확장합니다.
- 필요할 때의 exactly-once와 순서 있는 수집: Named Channels는 offset token으로 커밋된 진행 상황을 추적하고 각 채널 내에서 행 순서를 보존합니다. 소스 파티션에 자연스럽게 매핑되며 엄격한 exactly-once 복구를 간단하게 만듭니다.
- 높은 처리량, 낮은 지연 시간: 테이블당 최대 20 GB/s의 수집 속도와 최저 5초의 수집-조회 가능 지연 시간을 지원하도록 설계되었습니다. 결과는 워크로드 형태와 구성에 따라 달라집니다.
- 수집 중 변환(In-flight transformations): PIPE 객체 내에서 COPY 명령 구문으로 수집 중에 데이터를 정제·재구성·변환할 수 있습니다. 별도의 ETL 단계 없이 대상 테이블에 커밋되기 전에 컬럼 재정렬, 타입 캐스트, 표현식 적용을 할 수 있어요.
- 수집 시점 사전 클러스터링: 클러스터링 키가 있는 테이블에서 최적화된 쿼리 성능을 위해 수집 중에 데이터를 정렬합니다.
- Apache Iceberg 테이블 지원: Iceberg v2 및 Iceberg v3 테이블을 포함해 Snowflake 관리 Iceberg 테이블로 데이터를 스트리밍할 수 있습니다. 자세한 내용은 Apache Iceberg™ 테이블과 함께하는 Snowpipe Streaming 고성능 아키텍처를 참고하세요.
- 스키마 진화(Schema evolution): 변화하는 데이터 구조에 맞춰 테이블 스키마를 자동으로 적응시킵니다. Snowflake는 수동 DDL 변경 없이 들어오는 스트림에서 감지된 새 컬럼을 추가할 수 있어요.
- 모니터링과 오류 조사: 이벤트 테이블 텔레메트리를 조회해 수집 진행 상황, 처리 시간, 오류를 추적합니다. 오류 로깅을 활성화해 거부된 행 데이터를 검사하거나 재처리할 수 있어요.
- 간소화된 파이프라인: 전용 수집 인프라와 파이프라인 구성 요소가 줄어들어 운영 부담이 줄어듭니다.
연결 방법
Snowpipe Streaming은 다양한 워크로드에 맞는 여러 수집 경로를 지원합니다.
| 통합 | 가장 적합한 대상 |
|---|---|
| Java SDK (Java API 레퍼런스) | 고처리량 커스텀 애플리케이션. Java 11 이상 필요. |
| Python SDK (Python API 레퍼런스) | 데이터 엔지니어링과 Python 네이티브 워크플로우. Python 3.9 이상 필요. |
| Node.js SDK (Node.js API 레퍼런스) | JavaScript와 TypeScript 애플리케이션. Node.js 20 이상 필요. |
| REST API | 경량 워크로드, IoT 장치, 엣지 배포. |
| Snowflake Connector for Kafka | Apache Kafka 토픽 수집. |
Java, Python, Node.js SDK는 공유된 Rust 기반 클라이언트 코어를 사용합니다. 행이 도착하면 추가하세요. SDK는 시간·크기 임계값을 사용해 append를 자동으로 버퍼링·일괄 처리하고, 압축과 Snowflake 전송을 처리합니다. 직접 REST 클라이언트는 행을 NDJSON(newline-delimited JSON, 한 줄에 JSON 객체 하나)으로 묶고 압축을 직접 처리합니다.
참고 가능하면 REST API보다 Snowpipe Streaming SDK를 사용해 자동 일괄 처리와 더 단순한 통합을 누리세요. SDK가 환경에 적합하지 않을 때 직접 REST를 사용하세요.
시작하려면 Elastic Channels 또는 Named Channels를 선택한 뒤 해당 섹션의 SDK 또는 REST 튜토리얼을 따르세요. 나란히 비교와 사용 사례 지침은 채널 유형 선택을 참고하세요. PIPE 객체, 채널, offset token, 지원 데이터 타입에 대한 기술적 세부 사항은 핵심 개념을 참고하세요.
이런 경우에 추천
- 테이블당 최대 20 GB/s의 처리량이 필요한 고용량 스트리밍 워크로드
- 최저 5초의 수집-조회 가능 지연 시간이 필요한 실시간 분석과 대시보드
- SDK 또는 REST API를 통한 Elastic Channels를 사용하는 IoT, 텔레메트리, 분산 애플리케이션
- exactly-once 전달 보장이 있는 Named Channels를 사용하는 CDC 파이프라인
- Snowflake Connector for Kafka를 사용하는 Apache Kafka 토픽 수집
- 오픈 테이블 포맷 분석을 위한 Apache Iceberg 테이블로의 스트리밍
참고 SQL 네이티브 스트리밍을 찾고 있나요? 선언적 스트리밍 파이프라인은 다이나믹 테이블과 스트림을 Tasks와 함께 보세요.
Snowpipe Streaming vs Snowpipe
Snowpipe Streaming과 Snowpipe는 서로 보완 관계예요. 데이터가 애플리케이션, 장치, 서비스에서 행으로 도착하고 낮은 지연 시간의 데이터 가용성이 필요하면 Snowpipe Streaming을 사용하세요. 파이프라인이 이미 클라우드 스토리지에 파일을 만들고 일괄 처리 지향의 비교적 높은 지연 시간 로딩이 허용되면 Snowpipe를 사용하세요.