튜토리얼: Elastic Channels 시작하기
튜토리얼: Elastic Channels 시작하기 (REST API)
Snowpipe Streaming REST API와 Elastic Channel, cURL, JWT로 Snowflake에 데이터를 스트리밍하는 방법을 알려드릴게요. Elastic Channel REST 경로는 채널을 열거나 offset·continuation token을 관리할 필요가 없어서 새 REST 통합의 권장 시작점입니다.
출처: Snowflake 문서
본문
참고 SDK가 제공하는 향상된 처리량과 더 단순한 오류 처리를 누리기 위해 SDK 시작 가이드로 시작하는 것을 권장합니다. SDK가 적합하지 않은 경량 워크로드에 REST API를 사용하세요.
순서 있는 exactly-once 수입을 위한 Named Channel REST 경로는 튜토리얼: Snowpipe Streaming REST API 시작하기를 참고하세요.
사전 요구사항
- 키 페어 인증으로 구성된 Snowflake 사용자. 공개 키를 등록하세요:
ALTER USER MY_USER SET RSA_PUBLIC_KEY = '<your-public-key>';
필요한 권한은 액세스 제어를 참고하세요.
- Snowflake 데이터베이스, 스키마, 대상 테이블:
CREATE OR REPLACE DATABASE MY_DATABASE;
CREATE OR REPLACE SCHEMA MY_SCHEMA;
CREATE OR REPLACE TABLE MY_TABLE (
id NUMBER,
c1 NUMBER,
ts STRING
);
curl ,jq , SnowSQL 설치.- Snowflake 계정 식별자(형식 1:
myorg-account123 ). 자세한 내용은 계정 식별자를 참고하세요.
1단계: JWT 생성 및 환경 변수 설정
SnowSQL로 JWT를 생성하세요:
snowsql --private-key-path rsa_key.p8 --generate-jwt \
-a <ACCOUNT_IDENTIFIER> \
-u MY_USER
주의 JWT를 안전하게 저장하세요. 로그나 셸 기록에 노출하지 마세요.
튜토리얼용 환경 변수를 설정하세요:
export JWT_TOKEN="PASTE_YOUR_JWT_TOKEN_HERE"
export ACCOUNT="<ACCOUNT_IDENTIFIER>" # for example, ab12345
export USER="MY_USER"
export DB="MY_DATABASE"
export SCHEMA="MY_SCHEMA"
export TABLE="MY_TABLE"
export CONTROL_HOST="${ACCOUNT}.snowflakecomputing.com"
2단계: Ingest 호스트 발견 및 구성
Snowpipe Streaming REST는 두 개의 호스트네임을 사용합니다. 첫 번째는 Snowflake 계정 엔드포인트(
AWS PrivateLink, Azure Private Link, Google Cloud Private Service Connect를 사용한다면, SYSTEM$GET_PRIVATELINK_CONFIG가 반환한
중요 Snowflake 계정 이름에 밑줄(_)이 있으면, scoped token을 생성하기 전에
INGEST_HOST 에서 모든 밑줄을 대시(-)로 바꾸세요. 이후 모든 API 호출에 변환된 값(대시 사용)을 사용하세요. 예:my_account.region.ingest.snowflakecomputing.com 은my-account.region.ingest.snowflakecomputing.com 이 됩니다.
Ingest 호스트를 발견하세요:
export INGEST_HOST=$(curl -sS -X GET \
-H "Authorization: Bearer ${JWT_TOKEN}" \
-H "X-Snowflake-Authorization-Token-Type: KEYPAIR_JWT" \
"https://${CONTROL_HOST}/v2/streaming/hostname")
echo "Ingest Host: $INGEST_HOST"
프라이빗 연결을 사용한다면 기존 Snowflake 프라이빗 엔드포인트로 라우팅하는
Ingest 호스트용 scoped token을 얻으세요:
export SCOPED_TOKEN=$(curl -sS -X POST "https://$CONTROL_HOST/oauth/token" \
-H 'Content-Type: application/x-www-form-urlencoded' \
-H "Authorization: Bearer ${JWT_TOKEN}" \
-d "grant_type=urn:ietf:params:oauth:grant-type:jwt-bearer&scope=${INGEST_HOST}")
echo "Scoped token obtained"
3단계: 샘플 행 생성
NDJSON(newline-delimited JSON) 형식으로 행 배치를 만드세요. 요청이 재시도될 때 다운스트림에서 중복 제거할 수 있도록 각 행에 안정적인 이벤트 식별자(
export NOW_TS=$(date -u +"%Y-%m-%dT%H:%M:%SZ")
export REQUEST_ID=$(uuidgen | tr '[:upper:]' '[:lower:]')
cat <<EOF > rows.ndjson
{"id":1,"c1":$RANDOM,"ts":"$NOW_TS"}
{"id":2,"c1":$RANDOM,"ts":"$NOW_TS"}
EOF
4단계: Elastic Channel에 행 append
Elastic Channels에서만 사용할 수 있는 테이블 엔드포인트로 행을 보내세요. 첫 요청에서 Snowflake가 관리되는 기본 pipe를 만들거나 해석합니다. 모든 스트리밍 pipe에는 암시적
curl -sS -X POST \
-H "Authorization: Bearer ${SCOPED_TOKEN}" \
-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 .
성공적인 HTTP 200 응답이 내구성 승인입니다. Snowflake가 요청 페이로드를 내구성 있게 버퍼링했다는 뜻이에요. 행이 즉시 쿼리 가능하다는 의미는 아닙니다.
중요 Elastic 전달은 최소 한 번(at least once)입니다. 요청이 모호한 응답(네트워크 타임아웃, 응답 없음, 서버 측 5xx)을 반환하고 재시도하면 Snowflake가 원래 요청을 이미 수락했을 수 있어요. 같은 rowset의 재시도마다 같은
requestId 를 전달해 지원을 위한 서버 측 상관관계를 가능하게 하세요.retryCount 쿼리 매개변수는 Snowflake가 재시도를 식별하는 데 도움을 줍니다. 첫 시도에서0 으로 설정하고 재시도마다 증가시키세요.retryCount 가0 보다 크면 중복이 가능하다는 신호입니다. 페이로드에 안정적인 이벤트 식별자를 추가하고, 중복이 문제가 될 때 다운스트림에서 조정·중복 제거하세요.
또는 커스텀 PIPE용 pipe 엔드포인트를 사용할 수 있어요:
export PIPE="MY_TABLE-STREAMING"
curl -sS -X POST \
-H "Authorization: Bearer ${SCOPED_TOKEN}" \
-H "Content-Type: application/x-ndjson" \
"https://${INGEST_HOST}/v2/streaming/data/databases/$DB/schemas/$SCHEMA/pipes/$PIPE/channels/ELASTIC/rows?requestId=${REQUEST_ID}&retryCount=0" \
--data-binary @rows.ndjson | jq .
참고 직접 REST 클라이언트는 append가 내부적으로 결합되는 SDK 사용자와 달리 일괄 처리와 압축을 소유합니다. 프로덕션 REST 요청에서는 제한된 NDJSON 배치를 만들고 경과 시간 임계값 후 부분 배치를 보내세요. ZSTD 또는 Gzip 압축을 사용하세요. 페이로드가 일치하는 형식으로 압축되어 있을 때만
Content-Encoding: zstd 또는Content-Encoding: gzip 을 추가하세요. 각 Elastic 요청은 최대 4 MB의 페이로드 데이터(압축을 쓰면 압축 후 네트워크로 보내지는 페이로드 크기)를 담을 수 있습니다.
5단계: 데이터 확인
Snowflake가 데이터를 처리·구체화할 몇 초를 기다린 뒤 대상 테이블을 조회하세요:
SELECT * FROM MY_DATABASE.MY_SCHEMA.MY_TABLE;
행이 승인됐는데 대상 테이블에 안 보이면 오류 테이블을 확인하세요:
SELECT * FROM MY_DATABASE.MY_SCHEMA.MY_TABLE__ERRORS
LIMIT 20;
오류 테이블에 대한 자세한 내용은 고성능 아키텍처 Snowpipe Streaming의 오류 로깅을 참고하세요.
6단계: 정리 (선택)
rm -f rows.ndjson
unset JWT_TOKEN SCOPED_TOKEN ACCOUNT USER DB SCHEMA TABLE CONTROL_HOST INGEST_HOST NOW_TS REQUEST_ID PIPE
문제 해결
- HTTP 401 (Unauthorized): JWT가 유효하고 만료되지 않았는지 확인하세요. 필요하면 재생성하세요.
- HTTP 404 (Not Found): 데이터베이스, 스키마, 테이블 이름이 올바르고 Snowflake 계정에 존재하는지 확인하세요.
- HTTP 429 (Too Many Requests): 재시도 지연에 무작위 변동(jitter)이 있는 지수 백오프로 재시도하세요. 고정된 예약된 요청 비율을 가정하지 마세요.
- 승인 후 대상 테이블에 행이 없음: 구체화에 시간을 주고 오류 테이블을 확인하세요.
프라이빗 연결 문제 해결
CONTROL_HOST 가 해석되지 않음:privatelink-account-url 을 사용했는지, 프라이빗 DNS 존이 애플리케이션 런타임에 연결되어 있는지 확인하세요.Get Hostname 은 성공하는데INGEST_HOST 가 해석되지 않음: 반환된 ingest 호스트네임을 프라이빗 DNS에 추가하고 기존 Snowflake 프라이빗 엔드포인트로 라우팅하세요.- DNS는 커넥터 밖에서 동작하는데 안에서 실패함: 커넥터 컨테이너나 런타임에서 테스트하고, DNS 변경 후 음수 DNS 응답을 캐시하는 장기 실행 워커를 재시작하세요.
- TLS 호스트네임 불일치: 반환된 ingest 호스트네임을 URL 호스트네임으로 유지하세요. 원시 프라이빗 엔드포인트 IP에 연결하거나 TLS SNI·HTTP
Host 헤더를 재작성하지 마세요.
다음 단계
- 모범 사례: 일괄 처리, 압축, 안정적인 이벤트 ID, 정상 종료.
- 오류 처리: 재시도, 중복 위험, 요청 ID 상관관계.
- REST API 레퍼런스: 전체 Elastic 엔드포인트 사양.
- 제한 사항: 요청 크기, 전달 보장, SDK 버전 요구사항.