Elastic Channels 작업
Elastic Channels 작업
Elastic Channels의 내구성 승인, 플러시, 종료, 수입 검증, 모니터링을 다룹니다. 액세스 권한은 액세스 제어를 참고하세요.
출처: Snowflake 문서
본문
SDK로 내구성 승인 추적
이벤트가 도착하면 append하고 SDK가 내부적으로 일괄 처리하게 두세요. Future나 Promise를 반환하는 append에서는 제한된 Futures·Promises 세트를 보관하고, 애플리케이션 체크포인트의 모든 이벤트가 성공할 때까지 주기적으로 기다리세요. 콜백으로 추적되는 append에서는 null이 아닌 token과 성공·오류 핸들러를 사용해 같은 결과 세트를 추적하세요. 체크포인트는 이미 제출된 이벤트를 추적하지, 제출을 기다리는 배치를 추적하지 않아요.
미완료 이벤트와 바이트 한도로 메모리를 제한하고, 경과 시간 트리거로 저용량 스트림을 체크포인트하세요. 수입이 일시 중지되거나 끝날 때도 체크포인트하세요. 체크포인트 후에도 재구성 없이 수입이 재개될 수 있도록 클라이언트를 열어 두세요. 이후 이벤트가 승인됐더라도 미해결 이벤트를 지나 체크포인트를 전진시키지 마세요.
Elastic Channel 승인을 기다리는 것은 Named Channel 커밋된 offset을 추적하는 것과는 달라요. Elastic Channels에는
플러시와 정상 종료
프로듀서를 종료하기 전에 보류 중인 데이터를 플러시하고 in-flight append가 완료될 때까지 기다리세요. 이는 제출됐지만 아직 승인되지 않은 데이터를 잃을 위험을 줄여줘요.
Java:
// Stop accepting new source work, then initiate flush
channel.initiateFlush();
// Wait for all pending data to be flushed (with timeout)
CompletableFuture<Void> flush = channel.waitForFlush(Duration.ofMinutes(2));
flush.get(); // blocks until flush completes or throws on timeout/error
// Closing the client also closes the Elastic Channel
client.close();
Python:
# Stop accepting new source work
channel.initiate_flush()
# Wait for pending data to flush
channel.wait_for_flush(timeout_seconds=120)
client.close()
Node.js:
// Stop accepting new source work
channel.initiateFlush();
// Wait for pending data to flush
await channel.waitForFlush({ timeoutMs: 120000 });
await client.close();
클라이언트를 닫으면 Elastic Channel이 닫힙니다. 클라이언트가 닫힌 뒤에는 이후의 append 호출이 즉시 오류를 발생시킵니다(비동기 콜백이 아니라 동기적 closed-client 실패). REST 프로듀서의 경우 프로세스를 멈추기 전에 모든 in-flight 요청이 성공적인 HTTP 응답을 받았는지 확인하세요.
채널 상태
채널 상태를 조회해 Elastic Channel의 건강 상태를 확인할 수 있어요. Elastic Channel 이름은 항상
Java:
ChannelStatus status = channel.getChannelStatus();
System.out.println("Status: " + status.getStatusCode());
Python:
status = channel.get_channel_status()
print("Status:", status.status_code)
Node.js:
const status = await channel.getChannelStatus();
console.log("Status:", status.statusCode);
과거 채널 활동은 SNOWPIPE_STREAMING_CHANNEL_HISTORY 뷰를 참고하세요.
수입 검증
프로듀서를 실행한 후 데이터가 대상 테이블에 도달했는지 확인하세요. 구체화는 내구성 승인보다 수 초 뒤에 일어날 수 있어요.
SELECT COUNT(*) FROM MY_DATABASE.MY_SCHEMA.MY_TABLE;
SELECT * FROM MY_DATABASE.MY_SCHEMA.MY_TABLE
LIMIT 10;
오류 테이블 확인
내구성 승인 후 대상에서 행이 빠져 있다면 행 수준 처리 실패가 있는지 오류 테이블을 확인하세요.
SELECT * FROM MY_DATABASE.MY_SCHEMA.MY_TABLE__ERRORS
ORDER BY _SNOWFLAKE_INGEST_TIME DESC
LIMIT 20;
실패한 행을 포착하려면 대상 테이블에서 오류 로깅을 켜세요. 자세한 내용은 고성능 아키텍처 Snowpipe Streaming의 오류 로깅을 참고하세요.
수입 기록 확인
SNOWPIPE_STREAMING_CHANNEL_HISTORY 뷰는 모니터링과 문제 해결을 위한 채널 활동 기록을 제공합니다.
SELECT *
FROM SNOWFLAKE.ACCOUNT_USAGE.SNOWPIPE_STREAMING_CHANNEL_HISTORY
WHERE DATABASE_NAME = 'MY_DATABASE'
AND SCHEMA_NAME = 'MY_SCHEMA'
AND TABLE_NAME = 'MY_TABLE'
ORDER BY START_TIME DESC
LIMIT 20;
Prometheus 메트릭
메트릭 설정, Prometheus 구성, 클라이언트 로깅은 Prometheus와 로그로 SDK 클라이언트 모니터링을 참고하세요.