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로 마이그레이션하려면 다음 상위 수준 단계를 완료하세요.

  1. 각 대상 테이블에 대해 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];
  1. 모든 클래식 클라이언트에서 수집을 중지합니다.
  2. 클래식 클라이언트의 각 채널에 대해 마지막 커밋된 offset을 확인합니다. 이 offsets를 검색하려면 클래식 SDK의 getLatestCommittedOffsetTokens() 메서드를 사용하세요. 이 offsets가 클라이언트 측 레코드와 일치하는지 검증합니다.
  3. 애플리케이션 코드를 업데이트합니다.

프로젝트 의존성을 고성능 SDK(Java, Python, Node.js)로 전환합니다.

  1. 다음 API 및 구성 변경 섹션에 설명된 대로 API 호출을 업데이트합니다.
  2. Snowflake의 마지막 커밋된 offset을 사용해 테이블/PIPE당 하나의 클라이언트를 초기화합니다.
  3. 새 클라이언트가 구성되고 안정되면 수집을 재개합니다.

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')));

참고

이 방법은 기본 파이프 및 자동 스키마 진화 기능과 호환되지 않습니다.

더 알아보기 (Learn more)