Snowpipe Streaming 마이그레이션 가이드
Snowpipe Streaming 마이그레이션 가이드
이 가이드는 클래식 Snowpipe Java SDK에서 고성능 Snowpipe Streaming SDK로 마이그레이션하는 방법을 설명합니다. 여기서 논의되는 아키텍처 변경과 API 업데이트는 고성능 아키텍처가 세 언어 모두에서 사용 가능하므로 Python 및 Node.js SDK로의 마이그레이션에도 적용됩니다. 이 문서의 코드 예시는 Java지만, 핵심 마이그레이션 원칙은 언어 전반에 걸쳐 일관됩니다.
이 가이드는 Snowpipe Streaming Classic 애플리케이션이 사용하는 순서 보장 및 offset-token 복구 모델을 보존하므로 Named Channels를 사용합니다. 이러한 보장이 필요하지 않은 새 워크로드에 대해서는 Elastic Channels 를 별도로 평가하세요.
출처: Snowflake 문서
본문
주요 아키텍처 변경 사항
다음 표는 고성능 Snowpipe Streaming SDK의 가장 중요한 아키텍처 변경 사항을 요약합니다. SDK에 대한 자세한 비교는 Named Channel과 클래식 SDK 비교 를 참고하세요.
| 영역 | Classic (snowflake-ingest-java) | High-Performance (snowpipe-streaming SDK) |
|---|---|---|
| 진입점 | 데이터가 테이블로 직접 수집됩니다. | 데이터는 변환과 스키마 적용을 지원하는 PIPE 객체를 통해 수집됩니다. |
| SDK / Core | Java SDK만. | 공유 Rust core를 가진 여러 언어(Java, Python, Node.js)의 SDK. |
| API 이름 | insertRow/insertRows, openChannel(request) |
appendRow/appendRows, openChannel(channelName, offsetToken) |
| 오류 처리 | 클라이언트 측 검증이 수행됩니다. | 더 풍부한 오류 피드백이 있는 서버 측 검증이 제공됩니다. |
| 역압(Backpressure) 처리 | 스레드를 잠들게 하여 차단/무응답 상태를 초래합니다. | 오류를 반환하므로 호출자가 백오프/재시도 전략을 구현할 수 있습니다. |
| 클라이언트-테이블 매핑 | 단일 클라이언트 객체가 모든 테이블에 채널을 열 수 있었습니다. | 단일 클라이언트 객체는 이제 단일 pipe 객체에 독점적으로 묶입니다. |
| 과금 | 컴퓨팅과 클라이언트 수 기반. | 수집된 GB당 고정 요금. |
| 스키마 / 변환 | 클라이언트 측에서 관리. | PIPE 정의를 통해 서버 측에서 관리. |
마이그레이션 프로세스
애플리케이션을 고성능 SDK로 마이그레이션하려면 다음 상위 수준 단계를 완료하세요.
- 각 대상 테이블에 대해 PIPE를 생성 합니다.
CREATE PIPE my_pipe
AS COPY INTO my_table
FROM TABLE (DATA_SOURCE(TYPE => 'STREAMING'))
MATCH_BY_COLUMN_NAME = CASE_INSENSITIVE
[CLUSTER_AT_INGEST_TIME = TRUE];
- 모든 클래식 클라이언트에서 수집을 중지합니다.
- 클래식 클라이언트의 각 채널에 대해 마지막 커밋된 offset을 확인합니다. 이 offsets를 검색하려면 클래식 SDK의
getLatestCommittedOffsetTokens()메서드를 사용하세요. 이 offsets가 클라이언트 측 레코드와 일치하는지 검증합니다. - 애플리케이션 코드를 업데이트합니다.
프로젝트 의존성을 고성능 SDK(Java, Python, Node.js)로 전환합니다.
- 다음 API 및 구성 변경 섹션에 설명된 대로 API 호출을 업데이트합니다.
- Snowflake의 마지막 커밋된 offset을 사용해 테이블/PIPE당 하나의 클라이언트를 초기화합니다.
- 새 클라이언트가 구성되고 안정되면 수집을 재개합니다.
API 및 구성 변경 사항
마이그레이션 중에 API 호출과 구성 설정에 다음 변경 사항을 적용해야 합니다.
클라이언트 초기화
- Classic:
builder(name) - High-performance:
builder(name, db, schema, pipeName)
채널
- Classic:
openChannel(OpenChannelRequest) - High-performance:
openChannel(channelName, offsetToken)은 채널과 상태를 모두 반환합니다.
수집 메서드
- Classic:
insertRow/insertRows(...) - High-performance:
appendRow/appendRows(...)
Offset 추적
- 클래식 SDK의
getLatestCommittedOffsetTokens(channels)메서드는 제한된 가시성을 제공하고 오류 컨텍스트가 부족합니다. - 고성능 SDK는 여전히
getLatestCommittedOffsetTokens(...)를 지원하지만, 견고한 모니터링을 위해getChannelStatuses(...)를 사용하는 것이 좋습니다. 이 메서드는 다음 작업을 수행합니다.
오프셋이 예상대로 전진하고 있는지 확인합니다.
- 채널별 오류 개수와 상세 오류 정보를 반환합니다.
- 데이터 파이프라인의 사전 예방적(proactive) 모니터링과 문제 해결을 가능하게 합니다.
반정형 데이터 처리
고성능 SDK로 마이그레이션할 때 데이터가 리터럴 문자열로 저장되지 않도록 애플리케이션이 ARRAY와 VARIANT 열에 데이터를 제공하는 방식을 검토하세요.
동작 변경
v2에서 ARRAY 열에 직렬화된 문자열 리터럴 — 예: "[1, 2, 3]" — 을 전달하면 해당 문자열 리터럴을 포함하는 단일 요소 배열이 생성됩니다. 클래식 아키텍처 동작을 유지하려면 다음 옵션 중 하나를 선택하세요.
옵션 1: 네이티브 객체 전달 (권장)
appendRow를 호출하기 전에 클라이언트 애플리케이션을 업데이트해 JSON 문자열을 네이티브 객체로 역직렬화하세요.
- Java: 배열에는
java.util.List, 객체에는java.util.Map을 사용하세요. - Python: 네이티브
list및dict타입을 사용하세요. - Node.js: 네이티브
Array및Object타입을 사용하세요.
이점: 기본 파이프와 자동 스키마 진화와 호환됩니다.
옵션 2: 파이프 측 변환
PARSE_JSON 함수를 사용해 변환 로직으로 Pipe 객체를 명시적으로 정의하세요.
예시 SQL
CREATE PIPE my_pipe AS
COPY INTO my_table (my_array_col)
FROM (SELECT PARSE_JSON($1:my_array_col) FROM TABLE(DATA_SOURCE(TYPE => 'STREAMING')));
참고
이 방법은 기본 파이프 및 자동 스키마 진화 기능과 호환되지 않습니다.