Snowpipe Streaming REST API 엔드포인트

Snowpipe Streaming REST API 엔드포인트

참고

가능하면 자동 배치와 더 간단한 통합의 이점을 얻기 위해 REST API 대신 Snowpipe Streaming SDK를 사용하세요. 환경에 SDK가 적합하지 않을 때 직접 REST를 사용하세요.

Snowpipe Streaming REST API는 경량 워크로드를 위해 설계되었으며, Snowpipe Streaming SDK 없이 외부 애플리케이션과 통합하는 유연한 방법을 제공합니다.

이 레퍼런스는 두 가지 수집 모드를 모두 문서화합니다.

출처: Snowflake 문서

본문

요청 헤더

다음 요청 헤더는 Snowpipe Streaming REST API의 모든 엔드포인트에 적용됩니다.

헤더 설명
Authorization 인증 토큰
X-Snowflake-Authorization-Token-Type (optional) JWT/OAuth
Content-Encoding (optional) 페이로드의 압축 형식을 지정합니다. 지원: gzip, zstd.
User-Agent (optional) 클라이언트 애플리케이션을 식별합니다. 권장 형식: SnowpipeStreamingSDK/{version} ({platform}) {LANGUAGE}/{language_version} (app={partner-name}). 예: SnowpipeStreamingSDK/1.0.0 (Linux amd64) PYTHON/3.11.0 (app=MyPartnerApp).

직접 REST 클라이언트는 행을 newline-delimited JSON (NDJSON)으로 그룹화해야 하며, 한 줄당 하나의 JSON 객체와 자체 압축 처리가 필요합니다. Elastic과 Named Channel append 요청 모두 4 MB 페이로드 한도가 있습니다(압축을 사용한 경우 압축 이후 네트워크로 전송되는 페이로드 크기). 행을 배치하고 ZSTD 또는 Gzip 압축을 사용해 요청 오버헤드를 줄이세요. 전송 전에 배치 크기와 경과 시간을 제한하세요.

Get Hostname

Get Hostname은 Snowpipe Streaming REST API와 상호 작용하는 데 사용되는 호스트 이름을 반환합니다. 각 계정에는 고유한 호스트 이름이 있습니다.

참고

프라이빗 연결(AWS PrivateLink, Azure Private Link, Google Cloud Private Service Connect)의 경우 이 엔드포인트를 프라이빗 계정 URL을 통해 호출하고, scoped-token 교환과 스트리밍 작업에 사용하기 전에 반환된 ingest 호스트 이름에 대해 프라이빗 DNS를 구성하세요. 전체 단계는 REST API 튜토리얼의 ingest 호스트 발견 및 구성 을 참고하세요.

GET /v2/streaming/hostname

응답:

{
  "hostname": "string"
}

응답 필드 설명:

필드 타입 설명
Hostname String 계정의 호스트 이름.

Exchange Scoped Token

Exchange Scoped Token은 Snowpipe Streaming API 관련 서비스에만 접근하는 데 사용할 수 있는 보안 토큰을 반환합니다. 이것은 고객에게 보안 보호를 제공합니다.

POST /oauth/token

요청:

속성 필수 구성 요소 설명
content_type 예 헤더 “application/x-www-form-urlencoded”
grant_type 예 페이로드 “urn:ietf:params:oauth:grant-type:jwt-bearer”
scope 예 페이로드 계정의 호스트 이름.

응답:

{
  "token": "string"
}

응답 필드 설명:

필드 타입 설명
Token String scoped token.

Elastic Channel 엔드포인트

Elastic Channels는 별도의 open-channel 작업을 요구하지 않습니다. NDJSON 행 배치를 테이블 또는 파이프의 암시적 ELASTIC 채널로 직접 보내세요.

POST /v2/streaming/data/databases/{databaseName}/schemas/{schemaName}/tables/{tableName}/rows
POST /v2/streaming/data/databases/{databaseName}/schemas/{schemaName}/pipes/{pipeName}/channels/ELASTIC/rows

테이블 엔드포인트는 Elastic Channels에만 해당됩니다. 첫 번째 요청에서 Snowflake는 관리형 기본 파이프를 생성하거나 해결합니다.

속성 필수 구성 요소 설명
databaseName 예 URI 데이터베이스 이름, 대소문자 무시.
schemaName 예 URI 스키마 이름, 대소문자 무시.
tableName 또는 pipeName 예 URI 기본 파이프를 위한 대상 테이블, 또는 커스텀 파이프 이름.
rows 예 본문 NDJSON 행. 최대 Elastic 요청 페이로드는 4 MB입니다(압축을 사용한 경우 압축 이후 네트워크로 전송되는 페이로드 크기).
requestId 아니요 쿼리 파라미터 요청을 추적하는 UUID. 동일한 행 집합(행 배치)의 모든 재시도에 같은 값을 사용하고, 각각의 고유한 행 집합에 대해 새 UUID를 생성하세요.
retryCount 아니요 쿼리 파라미터 0에서 시작하는 재시도 시도 번호. 재시도마다 증가시키세요. 0보다 큰 값은 중복 행이 가능함을 신호하며 중복이 발생했음을 증명하지는 않습니다.

Elastic 요청에는 offsetToken, startOffsetToken, endOffsetToken, 또는 continuationToken이 포함되어서는 안 됩니다.

Elastic 응답 및 전달 의미론

성공적인 HTTP 200 응답은 내구성 있는 승인입니다. Snowflake가 요청 페이로드를 내구성 있게 버퍼링했습니다. 행이 대상 테이블에서 즉시 쿼리 가능함을 의미하지는 않습니다. 행 수준 처리 오류는 오류 로깅이 활성화된 경우 오류 테이블 에 영속화됩니다.

{
  "message": "OK"
}

Elastic Channels는 순서 보장 없이 at-least-once 전달을 제공합니다. 모호한 응답 후 재시도하면 중복 행이 생길 수 있습니다. 행 페이로드에 안정적인 이벤트 식별자를 포함하고, 중복이 중요할 때 다운스트림에서 조정하거나 중복 제거하세요.

Elastic append 예시

export REQUEST_ID=$(uuidgen)

curl -sS -X POST \
  -H "Authorization: Bearer ***" \
  -H "Content-Type: application/x-ndjson" \
  "https://${INGEST_HOST}/v2/streaming/data/databases/$DB/schemas/$SCHEMA/tables/$TABLE/rows?requestId=$REQUEST_ID&retryCount=0" \
  --data-binary @rows.ndjson | jq .

이 행 집합을 재시도할 때는 REQUEST_ID를 재사용하고 retryCount를 증가시키세요. 커스텀 파이프를 사용하려면 테이블 경로를 /pipes/$PIPE/channels/ELASTIC/rows로 바꾸세요.

Named Channel 엔드포인트

Named Channels는 명시적인 채널 라이프사이클, continuation token, 소스 offset token을 요구합니다. 다음 작업은 순서가 보장되고 정확히 한 번 수집을 지원합니다.

다음 다이어그램은 Named Channel 요청 흐름을 보여줍니다.

Named Channel 열기

Open Channel 작업은 파이프 또는 테이블에 대해 새 채널을 생성하거나 엽니다. 채널이 이미 존재하면 Snowflake는 채널의 클라이언트 시퀀서(client sequencer)를 올리고 마지막 커밋된 offset token을 반환합니다.

PUT /v2/streaming/databases/{databaseName}/schemas/{schemaName}/pipes/{pipeName}/channels/{channelName}

요청:

속성 필수 구성 요소 설명
databaseName 예 URI 데이터베이스 이름, 대소문자 무시.
schemaName 예 URI 스키마 이름, 대소문자 무시.
pipeName 예 URI 파이프 이름, 대소문자 무시.
channelName 예 URI 생성하거나 다시 여는 채널의 이름, 대소문자 무시.
offset_token 아니요 페이로드 채널을 열 때 offset token을 설정하는 데 사용되는 문자열.
fail_on_uncommitted_rows 아니요 페이로드 Boolean. true일 때 채널에 커밋되지 않은 진행 중(in-flight) 데이터가 있으면 서버가 HTTP 409 Conflict (ERR_CHANNEL_HAS_UNCOMMITTED_DATA)로 요청을 거부합니다. 그렇지 않으면 진행 중인 데이터는 조용히 버려지고 파이프의 같은 위치에 있는 채널을 느리게 만들 수 있습니다 (기본: false).
requestId 아니요 쿼리 파라미터 시스템을 통해 요청을 추적하는 데 사용되는 범용 고유 식별자(UUID).

응답:

{
  "next_continuation_token": "string",
  "channel_status": {
    "database_name": "string",
    "schema_name": "string",
    "pipe_name": "string",
    "channel_name": "string",
    "channel_status_code": "string",
    "last_committed_offset_token": "string",
    "created_on_ms": "long",
    "rows_inserted": "int",
    "rows_parsed": "int",
    "rows_error_count": "int",
    "last_error_offset_upper_bound": "string",
    "last_error_message": "string",
    "last_error_timestamp": "timestamp_utc",
    "snowflake_avg_processing_latency_ms": "int"
  }
}

응답 필드 설명:

필드 타입 설명
next_continuation_token String 이후의 Append Rows 요청에 사용해야 하는 API 관리 토큰. 토큰은 호출 시리즈를 연결해 연속적이고 순서가 맞는 데이터 스트림을 보장하고 exactly-once 전달을 위한 세션 상태를 유지합니다.
channel_status Object 채널에 대한 다음 상세 정보를 가진 중첩 객체:
database_name (String): 파이프가 있는 데이터베이스의 이름.
schema_name (String): 파이프가 있는 스키마의 이름.
pipe_name (String): 사용 중인 특정 파이프의 이름.
channel_name (String): 스트리밍 채널의 이름.
channel_status_code (String): 채널의 현재 상태를 나타내는 코드; 예: “ACTIVE”.
last_committed_offset_token (String): 마지막으로 성공적으로 커밋된 offset을 나타내는 토큰.
created_on_ms (Long): 채널이 생성된 시점의 밀리초 단위 타임스탬프.
rows_inserted (Int): 성공적으로 삽입된 총 행 수.
rows_parsed (Int): 파싱된 총 행 수.
rows_error_count (Int): 오류가 발생한 총 행 수.
last_error_offset_upper_bound (String): 마지막 오류가 발생한 offset의 상한을 나타내는 토큰.
last_error_message (String): 발생한 마지막 오류의 메시지.
last_error_timestamp (Long): 마지막 오류의 밀리초 단위 타임스탬프.
snowflake_avg_processing_latency_ms (Int): Snowflake의 평균 처리 지연 시간(밀리초).

Named Channel에 행 추가

Append Rows 작업은 주어진 채널에 행 배치를 삽입합니다.

POST /v2/streaming/data/databases/{databaseName}/schemas/{schemaName}/pipes/{pipeName}/channels/{channelName}/rows

요청:

속성 필수 구성 요소 설명
databaseName 예 URI 데이터베이스 이름, 대소문자 무시.
schemaName 예 URI 스키마 이름, 대소문자 무시.
pipeName 예 URI 파이프, 대소문자 무시.
channelName 예 URI 채널 이름, 대소문자 무시.
continuationToken 예 쿼리 파라미터 Snowflake의 continuation token, 클라이언트 및 행 시퀀서를 모두 캡슐화합니다.
startOffsetToken 아니요 쿼리 파라미터 배치의 첫 번째 행에 대한 offset token.
endOffsetToken 아니요 쿼리 파라미터 배치의 마지막 행에 대한 offset token.
rows 예 페이로드 NDJSON 형식으로 수집될 실제 데이터 페이로드. 이 속성의 최대 허용 크기는 4 MB입니다.
requestId 아니요 쿼리 파라미터 시스템을 통해 요청을 추적하는 데 사용되는 UUID.

참고

NDJSON 페이로드 내의 JSON 텍스트는 RFC 8259 표준을 엄격히 준수해야 합니다. 각 JSON 텍스트 뒤에는 줄 바꿈 문자 \n(0x0A)이 와야 합니다. 줄 바꿈 문자 앞에 캐리지 리턴 \r(0x0D)을 삽입할 수도 있습니다.

응답:

{
  "next_continuation_token": "string"
}

응답 필드 설명:

필드 타입 설명
next_continuation_token string 클라이언트 및 행 시퀀서를 모두 캡슐화하는 Snowflake의 다음 continuation token. 다음 배치를 삽입하는 데 사용해야 합니다.

Named Channel 삭제

Drop Channel 작업은 서버 측에서 채널과 그 메타데이터를 삭제합니다.

DELETE /v2/streaming/databases/{databaseName}/schemas/{schemaName}/pipes/{pipeName}/channels/{channelName}

요청:

속성 필수 구성 요소 설명
databaseName 예 URI 데이터베이스 이름, 대소문자 무시
schemaName 예 URI 스키마 이름, 대소문자 무시
pipeOrTableName 예 URI 파이프 또는 테이블 이름, 대소문자 무시
channelName 예 URI 채널 이름, 대소문자 무시
fail_on_uncommitted_rows 아니요 페이로드 Boolean. true일 때 채널에 커밋되지 않은 진행 중 데이터가 있으면 서버가 HTTP 409 Conflict (ERR_CHANNEL_HAS_UNCOMMITTED_DATA)로 요청을 거부합니다. 그렇지 않으면 진행 중인 데이터는 조용히 버려지고 파이프의 같은 위치에 있는 채널을 느리게 만들 수 있습니다 (기본: false).
requestId 아니요 쿼리 파라미터 시스템을 통해 요청을 추적하는 데 사용되는 UUID

응답:

이 작업은 HTTP 상태 코드 외에 특정 성공 응답이 없는 페이로드를 반환합니다.

Named Channel 상태 일괄 조회

Bulk Get Channel Status 작업은 특정 클라이언트 시퀀서에 대한 채널의 상태를 반환합니다.

POST /v2/streaming/databases/{databaseName}/schemas/{schemaName}/pipes/{pipeName}:bulk-channel-status

요청:

속성 필수 구성 요소 설명
databaseName 예 URI 데이터베이스 이름, 대소문자 무시
schemaName 예 URI 스키마 이름, 대소문자 무시
pipeName 예 URI 파이프 이름, 대소문자 무시
channel_names 예 페이로드 고객이 상태를 얻으려는 String 채널 이름의 배열; 이름은 대소문자를 구분합니다. 예: {"channel_names":["channel1", "channel2"]}.

응답:

{
  "channel_statuses": {
    "channel1": {
      "channel_status_code": "String",
      "last_committed_offset_token": "String",
      "database_name": "String",
      "schema_name": "String",
      "pipe_name": "String",
      "channel_name": "String",
      "rows_inserted": "int",
      "rows_parsed": "int",
      "rows_errors": "int",
      "last_error_offset_upper_bound": "String",
      "last_error_message": "String",
      "last_error_timestamp": "timestamp_utc",
      "snowflake_avg_processing_latency_ms": "int"
    },
    "channel2": {
      "comment": "same structure as channel1"
    }
    "comment": "potentially other channels"
  }
}

참고

요청된 채널이 서비스에서 발견되지 않으면 응답 페이로드는 channel_statuses 객체 내에 해당 채널에 대한 항목이 없습니다.

각 채널의 channel_statuses 필드 설명:

필드 타입 설명
channel_status_code String 채널의 상태를 나타냅니다.
last_committed_offset_token String 최신 커밋된 offset token.
database_name String 채널이 속한 데이터베이스의 이름.
schema_name String 채널이 속한 스키마의 이름.
pipe_name String 채널이 속한 파이프의 이름.
channel_name String 채널의 이름.
rows_inserted int 이 채널에 삽입된 모든 행의 개수.
rows_parsed int 파싱되었지만 이 채널에 반드시 삽입되지는 않은 모든 행의 개수.
rows_errors int 이 채널에 삽입될 때 오류가 발생해 거부된 모든 행의 개수.
last_error_offset_upper_bound String 수집 오류의 상한. 오류는 이 커밋된 offset token에 있거나 그 이전에 위치합니다.
last_error_message String 민감한 고객 데이터를 수정(redact)한 해당 채널의 최신 오류 코드에 대응하는 사람이 읽을 수 있는 메시지.
last_error_timestamp timestamp_utc 마지막 오류가 발생한 시점의 타임스탬프.
snowflake_avg_processing_latency_ms int 이 채널의 평균 엔드투엔드 처리 시간.

오류 응답 구조

Snowpipe Streaming REST API는 오류 응답에 대해 JSON 페이로드를 반환합니다. 이 구조는 자동 오류 처리와 사람의 분석 모두에 실행 가능한 정보를 제공합니다.

응답 페이로드는 다음 구조를 가집니다.

{
  "code": "...",
  "message": "..."
}

응답 필드

필드 타입 설명
Code String 안정적인 프로그래밍 오류 코드. 자동 오류 처리와 로깅에 사용할 수 있습니다. 예를 들어 애플리케이션 로직이 특정 코드를 확인해 사전 정의된 작업을 트리거할 수 있습니다.
Message String 오류를 설명하는 사람이 읽을 수 있는 메시지. 이 메시지는 변경될 수 있으므로 자동 파싱에 사용해서는 안 됩니다.

예시

다음 예시는 받을 수 있는 오류 응답을 보여줍니다.

{
  "code": "STALE_CONTINUATION_TOKEN_SEQUENCER",
  "message": "Channel sequencer in the continuation token is stale. Please reopen the channel"
}

이 예시는 오래된 채널 시퀀서가 있는 continuation token을 사용하려는 시도에 대한 응답을 보여줍니다. code는 명확하고 기계가 읽을 수 있는 오류 식별자를 제공하며, message는 사용자에게 도움이 되는 설명 텍스트를 제공합니다.

더 알아보기 (Learn more)