Elastic Channels 개요

Elastic Channels 개요

애플리케이션과 장치에서 Snowflake 테이블 또는 Snowflake 관리 Iceberg 테이블로 이벤트를 직접 스트리밍하세요. Elastic Channels는 수입만을 위해 Kafka나 다른 중간 매개체를 도입하지 않고도 Snowflake로 직접 스트리밍하게 해줘요. Elastic Channels가 수입 확장과 채널 라이프사이클을 관리하고, SDK가 append를 자동으로 일괄 처리하며 일시적 실패를 재시도합니다. 대부분의 새 Snowpipe Streaming 애플리케이션에 권장되는 시작점입니다.

출처: Snowflake 문서

본문

애플리케이션이 순서 있는 수입이나 exactly-once 복구를 요구하지 않을 때 Elastic Channels를 사용하세요. 그런 보장이 필요하면 Named Channels를 사용하세요. Exactly-once 복구는 재생을 위해 소스나 내구성 있는 애플리케이션 관리 스토리지에 레코드가 보관되어 있어야 합니다.

Elastic Channels 동작 방식

스트리밍 pipe에는 암시적 ELASTIC 채널이 하나 있습니다. 많은 프로듀서가 같은 Elastic Channel에 동시에 append할 수 있어요. Snowflake가 애플리케이션이 채널을 만들고, 이름짓고, 열고, 닫고, 복구할 필요 없이 append를 서버 관리 리소스에 분산합니다.

엔드투엔드 경로:

  • SDK 사용자는 getElasticChannel(Python: get_elastic_channel)로 핸들을 얻고 행이 도착하면 append합니다. SDK가 내부적으로 append를 자동 일괄 처리합니다. 직접 REST 호출자는 NDJSON(newline-delimited JSON, 한 줄에 JSON 객체 하나)을 사용해 Elastic REST 엔드포인트로 배치를 보냅니다.
  • Snowflake가 append를 내구성 있게 버퍼링하고 내구성 승인을 반환합니다. SDK append Future나 Promise가 성공적으로 완료되거나, 등록된 성공 핸들러가 결과를 보고하거나, 직접 REST 요청이 HTTP 200을 반환합니다. 이 시점에 데이터는 아직 쿼리할 준비가 되지 않았어요.
  • pipe가 서버 측에서 행을 처리합니다: 스키마를 검증하고, 구성된 변환이나 사전 클러스터링을 적용하며, 대상 테이블에 커밋합니다.
  • 대상 테이블이 커밋된 행을 받고, 테이블 처리가 끝난 뒤 쿼리 가능해집니다.

Snowflake가 append를 승인하면 프로듀서는 보관한 복사본을 해제할 수 있어요. 데이터는 Snowflake에 내구성 있게 존재하며, 처리와 쿼리 가용성은 그 뒤에 따릅니다. 오류 로깅이 활성화되면 행 수준 처리 오류는 오류 테이블에 영속화됩니다.

미승인 이벤트를 전달 요구사항에 따라 계속 사용할 수 있게 유지하세요. 장애 중에는 수입을 일시 중지하거나 들어오는 이벤트를 내구성 스토리지에 보관해 계속 수집하세요. SDK의 인메모리 버퍼는 프로세스 크래시를 견디지 못해요. 보존·복구 패턴은 미승인 데이터 보호를 참고하세요.

전달 의미론

Elastic Channels는 at-least-once 전달을 제공합니다. 순서는 보장되지 않아요. 프로듀서가 타임아웃, 프로세스 재시작, 또는 다른 모호한 실패 후 append를 재시도하면 이전 시도가 이미 수락됐을 수 있고 대상 테이블에 중복 행이 있을 수 있어요. 중복이 문제가 될 때 데이터 모델에 안정적인 이벤트 식별자를 포함하고, 다운스트림에서 조정·중복 제거하세요.

지원되는 대상

Elastic Channels는 다음에 쓸 수 있습니다:

  • 표준 Snowflake 테이블
  • 파티션된 테이블을 포함한 Snowflake 관리 Iceberg v2·v3 테이블

Iceberg 지원 세부 사항은 Apache Iceberg™ 테이블과 함께하는 Snowpipe Streaming 고성능 아키텍처를 참고하세요. 지원 데이터 타입과 자동 스키마 진화는 테이블 지원과 스키마를 참고하세요.

Append tokens

Append token은 Snowflake가 내구성 있게 승인한 메시지를 식별하는 데 도움을 줘요. 행을 append할 때 식별자를 제공하면 SDK가 성공·오류 핸들러에 그 식별자를 반환해 결과를 보낸 메시지와 매칭할 수 있게 해줍니다.

단일 행에는 이벤트 ID를, 단일 append로 제출한 행 그룹에는 그룹 식별자를 사용하세요. 성공 핸들러는 append가 내구성 있음을 보고하므로 그 메시지의 보관 복사본을 해제할 수 있어요. 오류 핸들러는 실패한 append를 식별해 애플리케이션이 처리하게 합니다.

예를 들어 Java API는 appendRow(Map<String, Object> row, Object appendToken)을 받아들입니다. 행을 보내기 전에 핸들러를 등록한 다음 token으로 "event-123"을 전달하세요:

channel.setSuccessHandler(detail ->
    System.out.println("Acknowledged: " + detail.getAppendTokens()));
channel.setErrorHandler(detail ->
    System.err.println("Failed: " + detail.getAppendTokens() + " cause: " + detail.getError()));
channel.appendRow(row, "event-123");

Snowflake가 이 행을 내구성 있게 승인하면 성공 핸들러가 detail.getAppendTokens()로 event-123을 받습니다. append가 실패하면 오류 핸들러가 token과 오류를 받아요. 콜백은 여러 append token을 함께 보고할 수 있습니다.

Java의 appendRowWithWait(row, appendToken)처럼 Future를 반환하는 append API를 사용한다면 반환된 Future로 완료를 추적할 수 있어요. SDK 튜토리얼이 Java, Python, Node.js 예제를 보여 줍니다.

Token을 작게 유지하세요. SDK는 append를 추적하는 동안 메모리에 보관하며 Snowflake로 보내지 않아요. Token은 결과 추적용 라벨이지, 중복 제거 키나 저장된 복구 위치가 아닙니다. 그 append에 콜백 보고가 필요 없다면 null 또는 None을 제공하세요.

직접 REST 요청 추적

직접 REST 호출자는 SDK append token을 사용하지 않습니다. requestId와 retryCount를 URL 쿼리 매개변수로 설정하세요:

POST /v2/streaming/data/databases/{databaseName}/schemas/{schemaName}/tables/{tableName}/rows?requestId={requestUuid}&retryCount=0

각 개별 rowset(행 배치)에 대해 새 request UUID를 생성하세요. 그 rowset을 재시도할 때는 UUID를 재사용하고 retryCount를 1, 2...로 증가시키세요. 이 매개변수들은 요청 상관관계를 지원합니다. 중복 제거 키는 아니에요. 자세한 내용은 Elastic REST 레퍼런스를 참고하세요.

성공·오류 콜백

setSuccessHandler와 setErrorHandler를 등록해 비동기 append 결과를 받으세요. 콜백은 공통 승인 또는 실패 이벤트로 그룹화된 append token을 받습니다.

콜백이 다른 승인을 지연시키지 않도록 빠르게 유지하세요. 재시도와 더 오래 걸리는 작업은 콜백 밖에서 처리하세요. 구현 지침은 콜백을 짧게 유지를 참고하세요.

시작하기

더 알아보기 (Learn more)