Snowpipe Streaming 모니터링
Snowpipe Streaming 모니터링
Snowpipe Streaming 파이프라인을 모니터링해 데이터가 도착하는지 확인하고 처리 지연을 식별하며 수입 실패를 조사하세요. Snowflake는 상황을 파악할 수 있게 해주는 두 가지 상호 보완적인 방법을 제공합니다:
- 이벤트 테이블로 수입 모니터링 — Snowflake 내부의 처리를 추적.
- Prometheus와 로그로 SDK 클라이언트 모니터링 — 애플리케이션에서 실행되는 SDK(software development kit)를 조사.
출처: Snowflake 문서
본문
이벤트 테이블로 수입 모니터링
이벤트 테이블을 사용해 Snowflake 안에서 스트리밍 파이프라인 전반의 처리를 모니터링하세요. 컬렉션을 활성화하면 Snowflake가 구성된 이벤트 테이블에 수입 활동을 기록합니다. SQL로 조회해 파이프라인을 비교하고, 문제를 조사하며, 애플리케이션에서 이러한 서버 측 이벤트를 수집하지 않고도 대시보드나 알림을 만들 수 있어요. 다음을 모니터링할 수 있습니다:
- 지연 시간(Latency): Snowflake가 데이터를 받은 때부터 대상 테이블에 커밋할 때까지의 시간. 재시도된 요청의 경우 측정은 가장 최근 수신부터 시작됩니다.
- 처리량과 진행 상황: 시간에 따라 수입된 행 수와, Named Channels의 경우 커밋과 함께 보고되는 offset tokens.
- 오류: 거부된 행 수, 오류 세부 정보, 채널 실패.
- 채널 활동: 채널 열기(OPEN)와 삭제(DROP) 작업.
사전 요구사항
Snowpipe Streaming 텔레메트리를 조회하기 전에 다음이 있는지 확인하세요:
- 고성능 아키텍처를 사용하는 Snowpipe Streaming 파이프라인.
- 텔레메트리를 받도록 구성된 대상인 활성 이벤트 테이블. Snowflake 관리 기본값은
SNOWFLAKE.TELEMETRY.EVENTS 입니다. - 그 테이블이나 그 뷰를 조회할 권한이 있는 역할. 기본 테이블의 경우
SNOWFLAKE.EVENTS_VIEWER 애플리케이션 역할이 기본 테이블이 아니라SNOWFLAKE.TELEMETRY.EVENTS_VIEW 에 대한 접근을 제공합니다. 뷰어 접근이 있다면 예제에서 뷰를 사용하세요. 다른 접근 옵션은 이벤트 테이블 액세스 제어를 참고하세요.
이벤트 수집 활성화
ALTER SCHEMA <database_name>.<schema_name> SET LOG_EVENT_LEVEL = INFO;
스키마가 수준을 설정하지 않으면 데이터베이스 설정을 상속하고, 이는 차례로 계정 설정을 상속합니다. 결과 설정이 수집되는 이벤트를 결정합니다:
INFO : 이 페이지에 설명된 다섯 가지 이벤트 유형을 모두 기록합니다.ERROR : 행·채널 오류만 기록합니다.OFF : Snowpipe Streaming 이벤트를 기록하지 않습니다. 수입은 정상적으로 계속됩니다.
변경 사항이 즉시 적용되지 않을 수 있어요. 자세한 내용은 텔레메트리 수준 설정을 참고하세요.
Snowpipe Streaming 이벤트 조회
각 이벤트는
SCOPE:"name" = 'snow.snowpipe.streaming' RECORD_TYPE = 'EVENT'
다음 쿼리는 하나의 대상 테이블에 대한 최근 15분 이벤트를 반환합니다. 이 페이지의 모든 예제에서
SELECT timestamp,
record:"name"::STRING AS event_name,
record:"severity_text"::STRING AS severity,
resource_attributes:"snow.database.name"::STRING AS database_name,
resource_attributes:"snow.schema.name"::STRING AS schema_name,
resource_attributes:"snow.table.name"::STRING AS table_name,
value:"channel_name"::STRING AS channel_name,
value
FROM <event_table>
WHERE scope:"name"::STRING = 'snow.snowpipe.streaming'
AND record_type = 'EVENT'
AND resource_attributes:"snow.database.name"::STRING = '<database_name>'
AND resource_attributes:"snow.schema.name"::STRING = '<schema_name>'
AND resource_attributes:"snow.table.name"::STRING = '<table_name>'
AND timestamp > DATEADD('minute', -15, CURRENT_TIMESTAMP())
ORDER BY timestamp DESC;
snow.executable.type : 서비스 식별자.SNOWPIPE_STREAMING 으로 설정됩니다.snow.database.name : 대상 테이블을 포함하는 데이터베이스.snow.schema.name : 대상 테이블을 포함하는 스키마.snow.table.name : 대상 테이블 이름.snow.pipe.name : 사용 가능할 때의 스트리밍 pipe 이름. 이 필드는 없을 수 있으므로 예제는 대신 데이터베이스, 스키마, 테이블로 필터링합니다. 채널 이름은VALUE:"channel_name" 에 있습니다.
이벤트 유형
모니터링하려는 것에 따라 이벤트 유형을 선택하세요. 심각도는 정보 이벤트(
| 이벤트 이름 | 심각도 | 용도 |
|---|---|---|
| commit | INFO | Snowflake가 처리된 데이터를 대상 테이블에 커밋할 때 행 수와 사용 가능한 offset tokens를 보고합니다. |
| latency | INFO | 채널에 대해 Snowflake가 데이터를 처리하는 데 걸린 시간을 보고합니다. |
| row_error | ERROR | 처리에 실패한 행에 대한 세부 정보를 제공합니다. |
| channel_lifecycle | INFO | 성공적인 채널 OPEN과 DROP 작업을 기록합니다. |
| channel_error | ERROR | 채널에 대한 작업이나 그 데이터 커밋 시 실패를 보고합니다. |
commit 이벤트로 수입 진행 상황 추적
channel_name : 커밋과 연관된 채널.row_count : 성공적으로 수입된 행 수.rows_parsed : 거부된 행을 포함해 파싱 중 읽힌 행 수.error_count : 검증 또는 파싱 중 거부된 행 수.uncompressed_bytes : 커밋 요청과 연관된 비압축 바이트. 요청에 여러 채널이 포함되면 같은 바이트 수가 각 채널의 이벤트에 나타날 수 있어요. 이 필드를 합산하면 데이터가 이중 계산될 수 있으므로 정확한 수입 또는 청구 합계에 사용하지 마세요.offset_start 와offset_end : 사용 가능할 때 커밋과 함께 보고되는 offset tokens. 이는 데이터를 보내는 애플리케이션이 정의한 문자열이지, Snowflake가 숫자나 알파벳 순서로 정렬하는 카운터가 아닙니다.
행 수는 각 이벤트를 설명하지, 누적 합계를 설명하지 않아요. Named Channels의 경우 채널 상태를 사용해 중단 후 수입을 재개할 위치를 결정하세요. 텔레메트리 이벤트 순서에 의존하지 마세요.
다음 쿼리는 대상 테이블과 채널별로 행·오류를 집계합니다. 오류율은 선택된 이벤트에서 거부된 행 수를 파싱된 행 수로 나눈 값이며, 파싱된 행이 없으면
SELECT resource_attributes:"snow.table.name"::STRING AS table_name,
value:"channel_name"::STRING AS channel_name,
SUM(value:"row_count"::NUMBER) AS rows_ingested,
SUM(value:"rows_parsed"::NUMBER) AS rows_parsed,
SUM(value:"error_count"::NUMBER) AS rows_with_errors,
SUM(value:"error_count"::NUMBER) / NULLIF(SUM(value:"rows_parsed"::NUMBER), 0) AS row_error_rate
FROM <event_table>
WHERE scope:"name"::STRING = 'snow.snowpipe.streaming'
AND record_type = 'EVENT'
AND record:"name"::STRING = 'commit'
AND resource_attributes:"snow.database.name"::STRING = '<database_name>'
AND resource_attributes:"snow.schema.name"::STRING = '<schema_name>'
AND resource_attributes:"snow.table.name"::STRING = '<table_name>'
AND timestamp > DATEADD('minute', -15, CURRENT_TIMESTAMP())
GROUP BY 1, 2
ORDER BY rows_with_errors DESC, rows_ingested DESC;
Named Channels의 경우 다음 쿼리는 각 채널에 대해 offset을 포함하는 최신 이벤트의
SELECT timestamp,
resource_attributes:"snow.table.name"::STRING AS table_name,
value:"channel_name"::STRING AS channel_name,
value:"offset_end"::STRING AS offset_end
FROM <event_table>
WHERE scope:"name"::STRING = 'snow.snowpipe.streaming'
AND record_type = 'EVENT'
AND record:"name"::STRING = 'commit'
AND resource_attributes:"snow.database.name"::STRING = '<database_name>'
AND resource_attributes:"snow.schema.name"::STRING = '<schema_name>'
AND resource_attributes:"snow.table.name"::STRING = '<table_name>'
AND value:"offset_end"::STRING IS NOT NULL
AND timestamp > DATEADD('minute', -15, CURRENT_TIMESTAMP())
QUALIFY RANK() OVER (PARTITION BY resource_attributes:"snow.table.name"::STRING,
value:"channel_name"::STRING ORDER BY timestamp DESC) = 1
ORDER BY channel_name;
latency 이벤트로 처리 시간 측정
Snowflake는 요청 타이밍이 활성화되고 사용 가능할 때 기록된 요청 시작 타임스탬프를 사용합니다. 그 외에는 측정이 Snowflake가 처리를 위해 데이터를 버퍼링할 때 시작됩니다.
이 측정은 요청이 Snowflake에 도달하기 전의 시간을 포함하지 않아요. 소스에서 레코드를 만들고 Snowflake에서 조회하는 총 시간이 아닙니다.
다음 쿼리는 측정값을 테이블·채널별로 1분 간격으로 그룹화합니다. 평균, 대략적 95번째 백분위수(p95), 최댓값, 측정 수를 반환합니다. p95 값은 측정의 약 95%가 그 값 이하에 해당하는 시간입니다:
SELECT DATE_TRUNC('minute', timestamp) AS minute,
resource_attributes:"snow.table.name"::STRING AS table_name,
value:"channel_name"::STRING AS channel_name,
AVG(value:"total_latency_ms"::NUMBER) AS avg_latency_ms,
APPROX_PERCENTILE(value:"total_latency_ms"::NUMBER, 0.95) AS p95_latency_ms,
MAX(value:"total_latency_ms"::NUMBER) AS max_latency_ms,
COUNT(value:"total_latency_ms"::NUMBER) AS measured_samples
FROM <event_table>
WHERE scope:"name"::STRING = 'snow.snowpipe.streaming'
AND record_type = 'EVENT'
AND record:"name"::STRING = 'latency'
AND resource_attributes:"snow.database.name"::STRING = '<database_name>'
AND resource_attributes:"snow.schema.name"::STRING = '<schema_name>'
AND resource_attributes:"snow.table.name"::STRING = '<table_name>'
AND timestamp > DATEADD('minute', -15, CURRENT_TIMESTAMP())
GROUP BY 1, 2, 3
ORDER BY minute DESC, p95_latency_ms DESC;
행·채널 오류 조사
channel_name : 오류가 발생한 채널.error_code 와error_message : 오류 식별자와 설명.error_col_name : 오류와 연관된 컬럼.error_line_number 와error_char_pos : 오류에 대해 보고된 줄 번호와 문자 위치.error_offset : 사용 가능할 때 오류와 연관된 offset token.
거부된 행마다 반드시
다음 쿼리는 최근 행·채널 오류를 반환합니다:
SELECT timestamp,
record:"name"::STRING AS event_name,
resource_attributes:"snow.table.name"::STRING AS table_name,
value:"channel_name"::STRING AS channel_name,
value:"error_type"::STRING AS error_type,
value:"error_code"::STRING AS error_code,
value:"error_message"::STRING AS error_message,
value:"error_offset"::STRING AS error_offset,
value:"error_col_name"::STRING AS error_column
FROM <event_table>
WHERE scope:"name"::STRING = 'snow.snowpipe.streaming'
AND record_type = 'EVENT'
AND record:"name"::STRING IN ('row_error', 'channel_error')
AND resource_attributes:"snow.database.name"::STRING = '<database_name>'
AND resource_attributes:"snow.schema.name"::STRING = '<schema_name>'
AND resource_attributes:"snow.table.name"::STRING = '<table_name>'
AND timestamp > DATEADD('minute', -15, CURRENT_TIMESTAMP())
ORDER BY timestamp DESC;
Snowflake가 영향을 받는 객체를 식별하지 못하면 이벤트에 객체 이름이 없을 수 있어요. 예제의 객체 이름 필터는 그런 이벤트를 제외합니다. 예상 오류가 누락됐으면 시간·이벤트 유형 필터는 유지하면서 그 필터를 제거하세요. 역할의 액세스 제어는 여전히 적용됩니다. 이벤트는 모든 클라이언트 인증·연결 실패를 포착하지는 않아요.
lifecycle 이벤트로 채널 활동 살펴보기
다음 쿼리는 최근 채널 작업을 반환합니다:
SELECT timestamp,
resource_attributes:"snow.table.name"::STRING AS table_name,
value:"channel_name"::STRING AS channel_name,
value:"channel_mode"::STRING AS channel_mode,
value:"event_type"::STRING AS event_type,
value:"tracking_label"::STRING AS tracking_label,
value:"client_version"::STRING AS client_version
FROM <event_table>
WHERE scope:"name"::STRING = 'snow.snowpipe.streaming'
AND record_type = 'EVENT'
AND record:"name"::STRING = 'channel_lifecycle'
AND resource_attributes:"snow.database.name"::STRING = '<database_name>'
AND resource_attributes:"snow.schema.name"::STRING = '<schema_name>'
AND resource_attributes:"snow.table.name"::STRING = '<table_name>'
AND timestamp > DATEADD('minute', -15, CURRENT_TIMESTAMP())
ORDER BY timestamp DESC;
Snowflake 알림 템플릿
알림 템플릿은 워크로드에 구성할 수 있는 사전 정의된 모니터링 조건을 제공합니다. Snowpipe Streaming 템플릿은 다음을 포함합니다:
SNOWPIPE_STREAMING_ROW_ERROR_RATE : 거부된 행과 파싱된 행의 비율을 모니터링합니다. 비율을 평가하기 전에 최소 파싱 행 수를 요구할 수 있습니다.SNOWPIPE_STREAMING_AUTH_FAILURE : 커밋 중에 보고된 선택된 권한, 객체 접근, 암호화 키 오류를 감지합니다. 모든 로그인·인증 토큰 실패를 다루지는 않아요.SNOWPIPE_STREAMING_CHANNEL_LIMIT : pipe가 채널 한도를 초과할 때(ERR_CHANNEL_LIMIT_EXCEEDED_FOR_PIPE ) 감지합니다. 데이터를 너무 빨리 보낼 때의 속도 제한은 감지하지 않아요.
계정, 데이터베이스, 스키마, 테이블을 모니터링할 수 있어요. 사용 가능한 템플릿과 설정은 계정별로 다를 수 있습니다. SYSTEM$LIST_ALERT_TEMPLATES로 사용 가능한 템플릿을 찾고, SYSTEM$GET_ALERT_TEMPLATE로 그 설정을 검사한 뒤 CREATE ALERT … FROM TEMPLATE로 알림을 만드세요. 자세한 내용은 Snowflake 알림을 참고하세요.
문제 해결
과거 채널 활동은 SNOWPIPE_STREAMING_CHANNEL_HISTORY 뷰를 사용하세요. 검사·재처리를 위해 거부된 행 데이터를 보존하려면 캡처·크기 한도를 적용받는 오류 테이블을 사용하세요.
- 이벤트가 나타나지 않음: 활성 이벤트 테이블, 조회 권한,
LOG_EVENT_LEVEL 설정을 확인하세요. 쿼리의 시간 범위와 객체 이름이 조사하려는 수입과 일치하는지 확인하세요. - 데이터가 도착하지 않음: 양수
row_count 를 가진 최근commit 이벤트를 찾고channel_error 이벤트를 확인하세요. Named Channels의 경우 채널 상태로 진행 상황을 확인하세요. - 행이 거부됨: commit 이벤트의
error_count 를rows_parsed 와 비교한 뒤row_error 이벤트에서 세부 사항을 조사하세요. - 처리가 느림: 같은 테이블, 채널, 시간 간격에 대해 지연 측정값을 행 수·오류와 비교하세요.
수집 동작과 비용
- 객체 이름 변경은 Snowflake가 이 정보를 잠시 캐시하므로 이벤트에 즉시 나타나지 않을 수 있어요.
- 특히 레코드가 많은 이벤트 테이블에서는 객체·시간 필터를 좁게 유지하세요. 수집, 저장, 조회 요금은 텔레메트리 데이터 수집 비용을 참고하세요.
Prometheus와 로그로 SDK 클라이언트 모니터링
SDK가 데이터를 보내는 애플리케이션에 대한 성능 메트릭과 로그를 제공합니다. 이는 이벤트 테이블 기록과 분리되어 있어요. Snowflake에서
Prometheus 메트릭 활성화·검증
Prometheus는 스크래핑(scraping)이라는 프로세스로 HTTP 엔드포인트를 주기적으로 읽어 메트릭을 수집합니다. SDK에는 애플리케이션이 실행되는 호스트에서 이 엔드포인트를 노출하는 메트릭 서버가 포함되어 있습니다.
- 애플리케이션이 시작되기 전에 환경에서
SS_ENABLE_METRICS 를true 로 설정하세요:
export SS_ENABLE_METRICS=true
- 그 환경에서 수입 애플리케이션을 시작하세요. 기본적으로 메트릭 서버는
127.0.0.1:50000 에서 수신하고/metrics 에서 메트릭을 제공합니다. - 애플리케이션이 실행되는 동안 같은 호스트의 다른 터미널에서 엔드포인트를 검증하세요:
curl http://127.0.0.1:50000/metrics
응답에는 Prometheus 텍스트 형식의 메트릭이 포함됩니다.
Prometheus 수집 구성
SDK 엔드포인트를 Prometheus 구성에 추가하세요. 이 예제는 Prometheus가 같은 호스트에서 실행되고 SDK의 기본 엔드포인트에 도달할 수 있다고 가정합니다:
scrape_configs:
- job_name: snowpipe_streaming_hp
metrics_path: /metrics
static_configs:
- targets: ['127.0.0.1:50000']
Prometheus가 다른 호스트나 별도 컨테이너에서 실행되면
클라이언트 로그 구성
클라이언트를 초기화하기 전에
클라이언트 문제 조사
같은 기간과 수입 대상을 대상으로 클라이언트 메트릭·로그를 이벤트 테이블 기록과 비교하세요. 채널별 상태와 복구는 Named Channel 작업 또는 Elastic Channel 작업을 참고하세요.
- 메트릭을 사용할 수 없음:
SS_ENABLE_METRICS 가 애플리케이션 시작 전에 설정됐는지, 애플리케이션이 여전히 실행 중인지 확인하세요. - 로컬 엔드포인트는 동작하는데 Prometheus가 메트릭을 수집하지 못함: 구성된 대상, 수신 주소, 포트, Prometheus와 애플리케이션 사이의 네트워크 접근을 확인하세요.
- 클라이언트 로그가 수입 오류를 설명하지 못함: Snowflake 내부 처리 중 오류에 대해 이벤트 테이블 기록을 확인하세요. 클라이언트 로그와 서버 측 이벤트는 수입의 서로 다른 부분을 다룹니다.