Snowpipe Streaming 고성능 아키텍처: 모범 사례

Snowpipe Streaming 고성능 아키텍처: 모범 사례

고성능 아키텍처의 Snowpipe Streaming을 사용해 견고한 데이터 수집 파이프라인을 설계하고 구현하기 위한 핵심 모범 사례를 안내하는 가이드입니다. 이 모범 사례들을 따르면 파이프라인이 내구성 있고 신뢰할 수 있으며 오류 처리가 효율적으로 이루어집니다.

출처: Snowflake 문서

본문

SDK가 자동으로 배치 처리하게 두기

행(rows)이 도착하는 대로 추가하세요. Java, Python, Node.js SDK는 내부적으로 시간과 크기 임계값을 기준으로 append를 버퍼링하고 결합하며, 압축과 Snowflake로의 데이터 전송을 처리합니다. 소스가 이미 여러 행을 함께 공급하는 경우에는 다중 행 append API를 사용하세요. 다중 행 append는 처리량을 위한 전제 조건이 아니라 활용하는 API입니다.

각 Named Channel에 대해 각 행의 커밋을 기다리지 말고 소스 순서대로 직렬로 행을 제출하세요. 진행 중인 작업(outstanding work)과 보관 중인 바이트를 제한하고, 주기적으로 커밋된 offset이 소스 체크포인트에 도달하거나 통과할 때까지 기다리세요. 커밋이 이루어진 후에만 소스 체크포인트를 전진시킵니다. 애플리케이션 체크포인트는 커밋된 진행 상황을 추적하는 것이지, SDK가 행을 보내기 위해 그룹화하는 방식을 추적하는 것이 아닙니다. 자세한 내용은 Named Channel 열고 사용하기 를 참고하세요. offset token 없는 Elastic acknowledgements에 대해서는 Elastic Channels 모범 사례 를 참고하세요.

Named Channel을 전략적으로 관리하기

성능과 장기적인 안정성을 위해 다음과 같은 채널 관리 전략을 적용하세요.

  • 장기 유지 채널 사용: 오버헤드를 최소화하려면 채널을 한 번 열고 수집 작업이 지속되는 동안 활성 상태로 유지하세요. 채널을 반복적으로 열고 닫는 것을 피하세요.
  • 결정적 채널 이름 사용: 예를 들어 source-env-region-client-id 같은 일관되고 예측 가능한 명명 규칙을 적용해 문제 해결을 간단히 하고 자동 복구 프로세스를 쉽게 하세요.
  • 여러 채널로 확장: 처리량을 늘리려면 여러 채널을 여세요. 서비스 한도와 처리량 요구 사항에 따라 이 채널들이 단일 대상 pipe 또는 여러 pipe를 가리킬 수 있습니다.
  • 채널 상태 모니터링: getChannelStatus 메서드를 정기적으로 사용해 수집 채널의 상태를 모니터링하세요.

last_committed_offset_token을 추적해 데이터가 제대로 수집되고 파이프라인이 진행되고 있는지 확인하세요. row_error_count를 모니터링해 잘못된 레코드나 다른 수집 문제를 조기에 감지하세요.

스키마를 일관되게 검증하기

들어오는 데이터가 예상 테이블 스키마를 준수하는지 확인해 수집 실패를 방지하고 데이터 무결성을 유지하세요.

  • 클라이언트 측 검증: 클라이언트 측에서 스키마 검증을 구현해 즉각적인 피드백을 제공하고 서버 측 오류를 줄이세요. 행 단위 전체 검증이 최대한의 안전성을 제공하지만, 성능이 더 좋은 방법으로는 배치 경계나 행 샘플링처럼 선택적 검증을 수행하는 것입니다.
  • 서버 측 검증: 고성능 아키텍처는 스키마 검증을 서버로 오프로드할 수 있습니다. 대상 pipe와 테이블로 수집하는 중에 스키마 불일치가 발생하면 오류와 그 개수가 getChannelStatus를 통해 보고됩니다.

클라이언트 측 메타데이터 열 추가하기

견고한 오류 감지와 복구를 활성화하려면 수집 메타데이터를 행 페이로드의 일부로 포함해야 합니다. 이를 위해서는 데이터 형태와 PIPE 정의를 미리 계획해야 합니다.

수집 전에 다음 열을 행 페이로드에 추가하세요.

  • CHANNEL_ID: 예를 들어 간결한 INTEGER
  • STREAM_OFFSET: 채널별로 단조 증가하는 BIGINT, 예를 들어 Kafka 파티션 offset

이 두 열은 함께 채널별 레코드를 고유하게 식별하며 데이터의 출처를 추적할 수 있게 합니다.

여러 pipe가 같은 대상 테이블로 수집하는 경우에는 선택적으로 PIPE_ID 열을 추가하세요. 이 열을 사용하면 행을 해당 수집 파이프라인까지 추적할 수 있습니다. 설명적인 pipe 이름은 별도의 조회 테이블에 저장하고, 압축된 정수로 매핑해 저장 비용을 줄일 수 있습니다.

메타데이터 offset을 사용해 오류 감지 및 복구하기

채널 모니터링과 메타데이터 열을 결합해 문제를 감지하고 복구하세요.

  • 상태 모니터링: getChannelStatus를 정기적으로 확인하세요. 증가하는 row_error_count는 잠재적 문제의 강한 신호입니다.
  • 누락 레코드 감지: 오류가 감지되면 SQL 쿼리를 사용해 STREAM_OFFSET 시퀀스의 공백을 확인함으로써 누락되거나 순서가 어긋난 레코드를 식별하세요.
SELECT
  PIPE_ID,
  CHANNEL_ID,
  STREAM_OFFSET,
  LAG(STREAM_OFFSET) OVER (
    PARTITION BY PIPE_ID, CHANNEL_ID
    ORDER BY STREAM_OFFSET
  ) AS previous_offset,
  (LAG(STREAM_OFFSET) OVER (
    PARTITION BY PIPE_ID, CHANNEL_ID
    ORDER BY STREAM_OFFSET
  ) + 1) AS expected_next
FROM my_table
QUALIFY STREAM_OFFSET != previous_offset + 1;

REST API 요청에 대해 행을 배치하고 압축 사용하기

직접 REST 클라이언트는 행을 스스로 그룹화하고 전송해야 합니다. 한 줄에 하나의 JSON 객체를 포함하는 newline-delimited JSON (NDJSON)을 사용해 행을 요청으로 결합하고, 압축을 사용해 네트워크 오버헤드를 줄이세요. 요청 크기를 제한하고 경과 시간에 따라 부분 배치를 전송해서 저용량 스트림이 전체 배치를 무한정 기다리지 않게 하세요.

Named Channel REST 요청은 압축을 사용한 경우 압축 이후 네트워크로 전송되는 페이로드에 4 MB 한도가 있습니다. 압축을 사용하면 각 요청에 더 많은 비압축 데이터 용량을 담을 수 있어 처리량을 높이고 필요한 API 호출 수를 줄일 수 있습니다.

Snowflake는 고성능 압축 알고리즘으로 ZSTD를 권장하지만 Gzip도 지원됩니다.

Snowflake 측 추적과 중복 감지를 위해 각 REST 요청에 requestId UUID 쿼리 파라미터를 포함하세요. 같은 행 집합(행 배치)에 대한 모든 재시도 시도에 동일한 requestId를 사용하세요. Snowflake는 요청 ID를 기록하고 이를 사용해 재시도된 요청을 연관시키고 잠재적 중복을 식별할 수 있습니다. 요청별 페이로드 한도는 제한 사항 및 고려 사항 을 참고하세요.

MATCH_BY_COLUMN_NAME으로 수집 성능과 비용 최적화하기

모든 데이터를 단일 VARIANT 열로 수집하는 대신, 소스 데이터에서 필요한 열을 매핑하도록 pipe를 구성하세요. 이를 위해 MATCH_BY_COLUMN_NAME = CASE_SENSITIVE를 사용하거나 pipe 정의에서 변환을 적용하세요. 이 모범 사례는 수집 비용을 최적화할 뿐만 아니라 스트리밍 데이터 파이프라인의 전반적인 성능도 향상시킵니다.

이 모범 사례는 다음과 같은 중요한 이점이 있습니다.

  • MATCH_BY_COLUMN_NAME = CASE_SENSITIVE를 사용하면 대상 테이블로 수집되는 데이터 값에 대해서만 요금이 부과됩니다. 반면 단일 VARIANT 열로 수집하면 키와 값 모두를 포함한 모든 JSON 바이트에 대해 요금이 부과됩니다. 장황하거나 많은 JSON 키가 있는 데이터의 경우 수집 비용이 크고 불필요하게 증가할 수 있습니다.
  • Snowflake의 처리 엔진이 계산적으로 더 효율적입니다. 전체 JSON 객체를 VARIANT로 파싱한 다음 필요한 열을 추출하는 대신, 이 방법은 필요한 값을 직접 추출합니다.

반정형 데이터에 네이티브 데이터 타입 사용하기

최적의 성능과 데이터 무결성을 위해 반정형 데이터를 직렬화된 문자열이 아닌 네이티브 언어 객체로 제공하세요.

  • 성능: 네이티브 객체를 사용하면 Snowflake 서버에서 추가 파싱 단계가 필요 없어 SDK가 데이터를 더 효율적으로 처리할 수 있습니다.
  • 타입 안전성: 고성능 아키텍처는 문자열 리터럴을 리터럴 텍스트로 취급합니다. 네이티브 객체를 사용하면 데이터가 이스케이프된 문자열 값이 아닌 구조화된 JSON으로 저장되도록 보장할 수 있습니다.
// 권장: SDK가 List를 구조화된 ARRAY로 변환합니다
row.put("tags", Arrays.asList("electronics", "sale"));
# 권장: SDK가 dict를 구조화된 VARIANT로 변환합니다
row["payload"] = {"event_id": 101, "status": "active"}
// 권장: ARRAY 열에는 네이티브 Array 사용
const row = {
  tags: ["electronics", "sale"],
  payload: { event_id: 101, status: "active" },
};
await channel.appendRow(row, "1");

Prometheus 지표 가져오기

지표 설정, Prometheus 구성, 클라이언트 로깅에 대해서는 Prometheus와 로그로 SDK 클라이언트 모니터링하기 를 참고하세요.

Linux와 macOS에서의 메모리 관리

SDK 버전 1.5.0부터 Snowpipe Streaming SDK는 지속적인 고처리량 스트리밍 워크로드에서 메모리 사용을 안정적으로 유지하기 위해 Linux와 macOS에서 jemalloc 메모리 할당자를 사용합니다. 이 변경은 Java, Python, Node.js SDK에 적용되며 자동으로 활성화됩니다. 애플리케이션을 업데이트하거나 구성을 설정할 필요가 없습니다.

이것이 애플리케이션에 의미하는 바는 다음과 같습니다.

  • 많은 채널을 열거나 높은 처리량으로 실행하는 장기 실행 수집 프로세스에서 더 예측 가능한 메모리 사용.
  • Linux 기본 시스템 할당자에 비해 메모리 단편화 감소.
  • SDK 공개 API나 기존 구성에는 변화가 없음.

Windows 빌드는 계속 시스템 할당자를 사용하며 이 변경의 적용을 받지 않습니다. Windows에서 SDK를 지속적으로 높은 처리량으로 실행하면서 메모리 사용이 꾸준히 증가하는 것을 관찰한다면, 워크로드를 Linux나 macOS에서 실행하거나 Snowflake 지원팀 에 문의하세요.

Java 배포를 위한 JVM 힙 크기 구성하기

SDK에는 JVM 힙 밖에서 메모리를 할당하는 네이티브 Rust 컴포넌트가 포함되어 있습니다. Java 래퍼를 사용할 때는 JVM 힙을 사용 가능한 메모리의 약 50%로 제한해 SDK의 네이티브 할당을 위한 공간을 남기세요.

예를 들어 8 GB RAM을 가진 호스트에서 -Xmx4g를 설정하세요:

MAVEN_OPTS="-Xmx4g" mvn exec:java -Dexec.mainClass="com.example.Main"
java -Xmx4g -jar your-app.jar

복원력을 위한 설계

수집을 try-catch 블록으로 감싸기

append 호출이 항상 성공한다고 가정하지 마세요. 동기 검증, 직렬화, 닫힌 클라이언트, 즉각적인 역압(backpressure) 실패는 append 호출에서 직접 발생합니다. SDK 오류(예: Python의 StreamingIngestError)를 잡아서 재시도하거나 다른 방식으로 처리하세요. 반환된 오류를 절대 무시하거나 버리지 마세요. HTTP 상태 코드를 해석하세요. 특히 Named Channel 무효화의 409와 제한(throttling)의 429에 주의하세요.

일부 오류는 행이 버퍼링된 후 비동기적으로 표면화되어 append 호출 시점에는 나타나지 않습니다. Named Channel 상태(예: getChannelStatus 또는 row_error_count)를 모니터링해 나중에 나타나는 실패를 감지하세요.

지수 백오프(exponential back-off) 구현하기

재시도 가능한 오류(429, 500, 503)에 대해서는 즉시 재시도하지 마세요. 시스템이 복구될 수 있도록 각 재시도 사이의 대기 시간을 늘리는 지수 백오프 전략을 사용하세요.

offset token으로 진행 상황 확인하기

getLatestCommittedOffsetToken을 주기적으로 호출해 어떤 데이터가 성공적으로 영속화되었는지 추적하세요. 409 오류가 발생하면 이 token을 사용해 채널을 다시 연 후 정확히 어느 지점부터 데이터를 재생해야 하는지 식별하세요.

채널 상태 모니터링하기

getChannelStatus()를 정기적으로 확인하세요. 상태 코드가 SUCCESS가 아니면 오류 처리 로직을 트리거해 채널 또는 클라이언트 연결을 리셋하세요.

더 알아보기 (Learn more)