Snowpipe Streaming 모니터링

Snowpipe Streaming 모니터링

Snowpipe Streaming 파이프라인을 모니터링해 데이터가 도착하는지 확인하고 처리 지연을 식별하며 수입 실패를 조사하세요. Snowflake는 상황을 파악할 수 있게 해주는 두 가지 상호 보완적인 방법을 제공합니다:

출처: 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에 대한 접근을 제공합니다. 뷰어 접근이 있다면 예제에서 뷰를 사용하세요. 다른 접근 옵션은 이벤트 테이블 액세스 제어를 참고하세요.

이벤트 수집 활성화

LOG_EVENT_LEVEL 매개변수는 Snowflake가 기록하는 이벤트를 제어합니다. 대상 테이블이 포함된 스키마에서 INFO로 설정하세요:

ALTER SCHEMA <database_name>.<schema_name> SET LOG_EVENT_LEVEL = INFO;

스키마가 수준을 설정하지 않으면 데이터베이스 설정을 상속하고, 이는 차례로 계정 설정을 상속합니다. 결과 설정이 수집되는 이벤트를 결정합니다:

  • INFO: 이 페이지에 설명된 다섯 가지 이벤트 유형을 모두 기록합니다.
  • ERROR: 행·채널 오류만 기록합니다.
  • OFF: Snowpipe Streaming 이벤트를 기록하지 않습니다. 수입은 정상적으로 계속됩니다.

변경 사항이 즉시 적용되지 않을 수 있어요. 자세한 내용은 텔레메트리 수준 설정을 참고하세요.

Snowpipe Streaming 이벤트 조회

각 이벤트는 RECORD에 이름과 심각도(정보 또는 오류), RESOURCE_ATTRIBUTES에 객체 이름, VALUE에 이벤트별 필드를 가집니다. SCOPE 컬럼은 이벤트를 생성한 서비스를 식별합니다. Snowpipe Streaming 이벤트를 선택하려면 두 필터를 모두 사용하세요:

  • SCOPE:"name" = 'snow.snowpipe.streaming'
  • RECORD_TYPE = 'EVENT'

다음 쿼리는 하나의 대상 테이블에 대한 최근 15분 이벤트를 반환합니다. 이 페이지의 모든 예제에서 <event_table>을 database.schema.object 형식의 이벤트 테이블 또는 뷰 이름으로 바꾸세요. 데이터베이스, 스키마, 테이블 자리표시자를 Snowflake에 저장된 정확한 이름(대문자 포함)으로 바꾸세요:

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;

RESOURCE_ATTRIBUTES의 필드가 영향을 받는 객체를 식별합니다:

  • snow.executable.type: 서비스 식별자. SNOWPIPE_STREAMING으로 설정됩니다.
  • snow.database.name: 대상 테이블을 포함하는 데이터베이스.
  • snow.schema.name: 대상 테이블을 포함하는 스키마.
  • snow.table.name: 대상 테이블 이름.
  • snow.pipe.name: 사용 가능할 때의 스트리밍 pipe 이름. 이 필드는 없을 수 있으므로 예제는 대신 데이터베이스, 스키마, 테이블로 필터링합니다. 채널 이름은 VALUE:"channel_name"에 있습니다.

이벤트 유형

모니터링하려는 것에 따라 이벤트 유형을 선택하세요. 심각도는 정보 이벤트(INFO)와 오류(ERROR)를 식별합니다:

이벤트 이름 심각도 용도
commit INFO Snowflake가 처리된 데이터를 대상 테이블에 커밋할 때 행 수와 사용 가능한 offset tokens를 보고합니다.
latency INFO 채널에 대해 Snowflake가 데이터를 처리하는 데 걸린 시간을 보고합니다.
row_error ERROR 처리에 실패한 행에 대한 세부 정보를 제공합니다.
channel_lifecycle INFO 성공적인 채널 OPEN과 DROP 작업을 기록합니다.
channel_error ERROR 채널에 대한 작업이나 그 데이터 커밋 시 실패를 보고합니다.

commit 이벤트로 수입 진행 상황 추적

commit 이벤트의 VALUE 컬럼에는 다음 필드가 포함됩니다:

  • channel_name: 커밋과 연관된 채널.
  • row_count: 성공적으로 수입된 행 수.
  • rows_parsed: 거부된 행을 포함해 파싱 중 읽힌 행 수.
  • error_count: 검증 또는 파싱 중 거부된 행 수.
  • uncompressed_bytes: 커밋 요청과 연관된 비압축 바이트. 요청에 여러 채널이 포함되면 같은 바이트 수가 각 채널의 이벤트에 나타날 수 있어요. 이 필드를 합산하면 데이터가 이중 계산될 수 있으므로 정확한 수입 또는 청구 합계에 사용하지 마세요.
  • offset_start와 offset_end: 사용 가능할 때 커밋과 함께 보고되는 offset tokens. 이는 데이터를 보내는 애플리케이션이 정의한 문자열이지, Snowflake가 숫자나 알파벳 순서로 정렬하는 카운터가 아닙니다.

행 수는 각 이벤트를 설명하지, 누적 합계를 설명하지 않아요. Named Channels의 경우 채널 상태를 사용해 중단 후 수입을 재개할 위치를 결정하세요. 텔레메트리 이벤트 순서에 의존하지 마세요.

다음 쿼리는 대상 테이블과 채널별로 행·오류를 집계합니다. 오류율은 선택된 이벤트에서 거부된 행 수를 파싱된 행 수로 나눈 값이며, 파싱된 행이 없으면 NULL입니다:

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을 포함하는 최신 이벤트의 offset_end를 반환합니다. 여러 이벤트가 최신 타임스탬프를 공유하면 쿼리가 모두 반환합니다. 이는 소스의 레코드 순서가 아니라 최근에 기록된 이벤트를 식별합니다. 배경은 Offset tokens와 exactly-once 전달을 참고하세요:

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 이벤트로 처리 시간 측정

latency 이벤트에는 channel_name과, 측정이 가능할 때 total_latency_ms(밀리초 단위 처리 시간)가 포함됩니다. 측정은 데이터가 Snowflake의 수입 엔드포인트에 도달할 때 시작해 Snowflake가 대상 테이블에 커밋할 때 끝납니다. 재시도된 요청의 경우 시작 시간은 첫 시도가 아니라 그 엔드포인트의 가장 최근 수신입니다.

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;

행·채널 오류 조사

row_error 이벤트는 다음 필드를 포함할 수 있습니다:

  • channel_name: 오류가 발생한 채널.
  • error_code와 error_message: 오류 식별자와 설명.
  • error_col_name: 오류와 연관된 컬럼.
  • error_line_number와 error_char_pos: 오류에 대해 보고된 줄 번호와 문자 위치.
  • error_offset: 사용 가능할 때 오류와 연관된 offset token.

channel_error 이벤트는 채널 이름, 오류 코드, 오류 메시지, error_type 분류를 포함합니다. 커밋 중의 실패는 error_type이 SYSTEM으로 설정되고 uncompressed_bytes를 포함할 수 있어요. 다른 채널 작업 실패는 USER 또는 SYSTEM으로 분류됩니다. 오류 코드와 메시지를 사용해 취할 조치를 결정하세요. 분류 자체는 원인을 확립하지 않습니다.

거부된 행마다 반드시 row_error 이벤트가 있는 건 아니에요. 합계에는 commit 이벤트의 error_count를 사용하고, 개별 실패 조사에는 row-error 이벤트를 사용하세요.

다음 쿼리는 최근 행·채널 오류를 반환합니다:

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 이벤트로 채널 활동 살펴보기

channel_lifecycle 이벤트는 성공적인 OPEN 또는 DROP 작업을 기록합니다. 그 VALUE 필드는 채널(channel_name), 모드(channel_mode), 작업(event_type)을 식별합니다. 클라이언트가 제공하면 tracking_label과 client_version이 작업을 수행한 클라이언트에 대한 추가 정보를 제공합니다.

다음 쿼리는 최근 채널 작업을 반환합니다:

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에서 LOG_EVENT_LEVEL을 활성화해도 SDK 메트릭이 켜지거나 SDK 로그 수준이 바뀌지 않습니다.

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 텍스트 형식의 메트릭이 포함됩니다. SS_METRICS_IP와 SS_METRICS_PORT 환경 변수가 수신 주소와 포트를 제어합니다. 이들의 기본값과 다른 설정은 환경 변수를 참고하세요.

Prometheus 수집 구성

SDK 엔드포인트를 Prometheus 구성에 추가하세요. 이 예제는 Prometheus가 같은 호스트에서 실행되고 SDK의 기본 엔드포인트에 도달할 수 있다고 가정합니다:

scrape_configs:
  - job_name: snowpipe_streaming_hp
    metrics_path: /metrics
    static_configs:
      - targets: ['127.0.0.1:50000']

Prometheus가 다른 호스트나 별도 컨테이너에서 실행되면 127.0.0.1은 수입 애플리케이션이 아니라 그 호스트·컨테이너를 가리킵니다. SDK 수신 주소와 Prometheus 대상이 통신할 수 있게 구성하세요. 메트릭 엔드포인트에 대한 접근을 모니터링 인프라로 제한하세요. 공개적으로 노출하지 마세요.

클라이언트 로그 구성

클라이언트를 초기화하기 전에 SS_LOG_LEVEL을 설정해 로그 출력을 제어하세요. 지원되는 값은 info, warn, error이며 기본값은 info입니다. 로깅·메트릭 설정은 같은 프로세스의 모든 SDK 클라이언트에 적용됩니다. 구성 레퍼런스는 환경 변수를 참고하세요.

클라이언트 문제 조사

같은 기간과 수입 대상을 대상으로 클라이언트 메트릭·로그를 이벤트 테이블 기록과 비교하세요. 채널별 상태와 복구는 Named Channel 작업 또는 Elastic Channel 작업을 참고하세요.

  • 메트릭을 사용할 수 없음: SS_ENABLE_METRICS가 애플리케이션 시작 전에 설정됐는지, 애플리케이션이 여전히 실행 중인지 확인하세요.
  • 로컬 엔드포인트는 동작하는데 Prometheus가 메트릭을 수집하지 못함: 구성된 대상, 수신 주소, 포트, Prometheus와 애플리케이션 사이의 네트워크 접근을 확인하세요.
  • 클라이언트 로그가 수입 오류를 설명하지 못함: Snowflake 내부 처리 중 오류에 대해 이벤트 테이블 기록을 확인하세요. 클라이언트 로그와 서버 측 이벤트는 수입의 서로 다른 부분을 다룹니다.

더 알아보기 (Learn more)