Snowpipe Streaming Classic과 함께하는 Snowflake Connector for Kafka
Snowpipe Streaming Classic과 함께하는 Snowflake Connector for Kafka
Kafka 토픽의 데이터를 Snowpipe Streaming으로 Snowflake 테이블에 로드하는 방법을 알려드릴게요. 지정된 플러시 버퍼 임계값(시간, 메모리, 메시지 수)에 도달하면 커넥터가 Snowpipe Streaming API를 호출해 데이터 행을 Snowflake 테이블에 씁니다. 이 아키텍처는 비슷한 데이터 볼륨을 로드할 때 더 낮은 로드 지연 시간과 그에 따른 더 낮은 비용을 제공합니다.
출처: Snowflake 문서
본문
중요
- 사전 예고: 이 페이지는 Kafka 커넥터(v3 이하)와 함께 Snowpipe Streaming classic을 사용하는 방법을 문서화합니다. 새 구현에서는 Snowflake Connector for Kafka (v4)를 권장하며, 이는 기본적으로 고성능 아키텍처를 사용합니다.
- 즉시 변경할 필요는 없습니다. 현재 워크로드는 완전히 지원됩니다.
Kafka에서의 데이터 로딩 체인에서 Snowpipe를 Snowpipe Streaming으로 대체할 수 있어요. 지정된 플러시 버퍼 임계값(시간, 메모리, 메시지 수)에 도달하면 커넥터가 Snowpipe Streaming API("API")를 호출해 데이터 행을 Snowflake 테이블에 씁니다. 이 아키텍처는 비슷한 데이터 볼륨을 로드할 때 더 낮은 로드 지연 시간과 그에 따른 더 낮은 비용을 제공합니다.
Snowpipe Streaming Classic과 함께 사용하려면 Kafka 커넥터 버전 2.0.0 이상이 필요합니다. Snowpipe Streaming Classic과 함께하는 Kafka 커넥터는 Snowflake Ingest SDK를 포함하며 Apache Kafka 토픽에서 대상 테이블로 행을 직접 스트리밍하는 것을 지원합니다.
최소 필요 버전
Snowpipe Streaming을 지원하는 Kafka 커넥터의 최소 버전은 2.0.0입니다.
Kafka 구성 속성
연결 설정을 Kafka 커넥터 속성 파일에 저장하세요. 자세한 내용은 Kafka 커넥터 구성을 참고하세요.
필수 속성
Kafka 커넥터 속성 파일에서 연결 설정을 추가하거나 편집하세요. 자세한 내용은 Kafka 커넥터 구성을 참고하세요.
SNOWPIPE_STREAMING SNOWPIPE (기본값)
토픽 데이터를 큐에 넣고 로드할 백엔드 서비스를 선택하는 데 필요한 추가 설정은 없어요. 평소처럼 Kafka 커넥터 속성 파일에 추가 속성을 구성하세요.
클라이언트 최적화 속성
값:
고처리량 시나리오(예: 커넥터당 50 MB/s)에서는 이 속성을 활성화하면 지연 시간이나 비용이 더 높아질 수 있습니다. 고처리량 시나리오에서는 이 속성을 비활성화하는 것을 권장합니다.
버퍼와 폴링 속성
값: 최솟값
Snowpipe Streaming은 Kafka 커넥터의 버퍼 플러시 시간과는 별개로 매 1초마다 자동으로 데이터를 플러시한다는 점에 유의하세요. Kafka 버퍼 플러시 시간에 도달한 뒤 데이터는 Snowpipe Streaming을 통해 1초의 지연 시간으로 Snowflake에 전송됩니다. 자세한 내용은 Snowpipe Streaming 지연 시간을 참고하세요.
값: 최솟값
값: 최솟값
값:
Kafka 커넥터 속성 외에도 Kafka 소비자
오류 처리와 DLQ 속성
NONE : 첫 오류를 만나면 데이터 로드를 중지합니다.ALL : 모든 오류를 무시하고 데이터 로드를 계속합니다.
기본값:
TRUE : 오류 메시지를 씁니다.FALSE : 오류 메시지를 쓰지 않습니다.
기본값:
값: 커스텀 텍스트 문자열. 기본값: 없음.
Exactly-once 의미론
Exactly-once 의미론은 중복이나 데이터 유실 없이 Kafka 메시지를 전달합니다. 이 전달 보장은 기본적으로 Snowpipe Streaming과 함께하는 Kafka 커넥터에 설정됩니다.
Kafka 커넥터는 파티션과 채널 사이에 일대일 매핑을 채택하고 두 가지 서로 다른 offset을 사용합니다.
- 소비자 offset(Consumer offset): 소비자가 소비한 가장 최근 offset을 추적하며 Kafka가 관리합니다.
- Offset token: Snowflake에서 커밋된 가장 최근 offset을 추적하며 Snowflake가 관리합니다.
Kafka 커넥터는 항상 누락된 offset을 처리하지는 않는다는 점에 유의하세요. Snowflake는 모든 레코드가 순차적으로 증가하는 offset을 가질 것으로 기대합니다. 누락된 offset은 특정 사용 사례에서 Kafka 커넥터를 망가뜨립니다. NULL 레코드 대신 tombstone 레코드를 사용하는 것을 권장합니다.
Kafka 커넥터는 다음 모범 사례를 구현해 exactly-once 전달을 달성합니다.
채널 열기/다시 열기:
- 주어진 파티션에 대한 채널을 열거나 다시 열 때, Kafka 커넥터는
getLatestCommittedOffsetToken API를 통해 Snowflake에서 검색한 최신 커밋된 offset token을 소스 오브 트루스로 사용하고 그에 따라 Kafka의 소비자 offset을 리셋합니다. - 소비자 offset이 더 이상 데이터 보존 기간 안에 없으면 예외가 발생하고 적절한 조치를 결정할 수 있어요.
- Kafka 커넥터가 Kafka의 소비자 offset을 리셋하지 않고 소스 오브 트루스로 사용하는 유일한 시나리오는 Snowflake의 offset token이 NULL일 때입니다. 이 경우 커넥터는 Kafka가 보낸 offset을 받아들이고 offset token은 이후에 업데이트됩니다.
레코드 처리:
- Kafka의 잠재적 버그에서 생길 수 있는 비연속 offset에 대한 추가 안전 계층을 보장하기 위해, Snowflake는 최신 처리된 offset을 추적하는 인메모리 변수를 유지합니다. Snowflake는 현재 행의 offset이 최신 처리된 offset + 1과 같을 때만 행을 수락하며, 이로써 수입 프로세스가 연속적이고 정확하도록 하는 추가 보호 계층을 더합니다.
예외, 실패, 크래시 복구 처리:
- 복구 프로세스의 일부로 Snowflake는 채널을 다시 열고 최신 커밋된 offset token으로 소비자 offset을 리셋하는 앞서 설명한 채널 열기/다시 열기 로직을 일관되게 따릅니다. 이렇게 하면 Snowflake가 Kafka에 최신 커밋된 offset token보다 하나 큰 offset부터 데이터를 보내라고 신호를 보내, 데이터 유실 없이 실패 지점에서 수입을 재개할 수 있습니다.
재시도 메커니즘 구현:
- 잠재적인 일시적 문제를 고려해 Snowflake는 API 호출에 재시도 메커니즘을 통합합니다. Snowflake는 이 API 호출을 여러 번 재시도해 성공 확률을 높이고 간헐적 실패가 수입 프로세스에 영향을 주는 위험을 완화합니다.
소비자 offset 전진:
- 정기적인 간격으로 Snowflake는 최신 커밋된 offset token을 사용해 소비자 offset을 전진시켜 수입 프로세스가 Snowflake의 최신 데이터 상태와 지속적으로 정렬되도록 합니다.
컨버터
Snowpipe Streaming은 다음과 같은 많은 커뮤니티 기반 컨버터를 지원합니다.
io.confluent.connect.avro.AvroConverter org.apache.kafka.connect.json.JsonConverter io.confluent.connect.protobuf.ProtobufConverter io.confluent.connect.json.JsonSchemaConverter org.apache.kafka.connect.converters.ByteArrayConverter org.apache.kafka.connect.storage.StringConverter
다른 커뮤니티 기반 컨버터도 지원될 수 있지만 검증되지는 않았습니다. Snowflake 컨버터는 Snowpipe Streaming에서 지원되지 않습니다.
Dead-letter queues
Snowpipe Streaming과 함께하는 Kafka 커넥터는 손상된 레코드나 실패로 인해 성공적으로 처리할 수 없는 레코드를 위한 DLQ(dead-letter queue)를 지원합니다. 모니터링에 대한 자세한 내용은 Apache Kafka 문서를 참고하세요.
스키마 감지와 스키마 진화
Snowpipe Streaming과 함께하는 Kafka 커넥터는 스키마 감지와 진화를 지원합니다. Snowflake 테이블의 구조를 정의·진화시켜 Kafka 커넥터가 로드하는 새 Snowpipe Streaming 데이터의 구조를 자동으로 지원할 수 있어요.
Snowpipe Streaming과 함께하는 Kafka 커넥터에 스키마 감지·진화를 활성화하려면 다음 Kafka 속성을 구성하세요.
snowflake.ingestion.method snowflake.enable.schematization schema.registry.url
자세한 내용은 Snowpipe Streaming Classic과 함께하는 Kafka 커넥터의 스키마 감지·진화를 참고하세요.
수입 지연 시간 추정
수입 지연 시간을 추정하려면 RECORD_METADATA의
참고 이 필드는 구성한 Snowpipe Streaming 지연 시간을 고려하지 않으므로 레코드가 Snowflake 테이블에 표시된 시점을 나타내지 않아요.
청구와 사용량
Snowpipe Streaming 청구 정보는 Snowpipe Streaming Classic 비용을 참고하세요.
제한 사항
Snowpipe Streaming 제한 사항
Snowpipe Streaming 제한 사항을 참고하세요.
페일오버 제한 사항
보조 페일오버 그룹이 주 그룹으로 승격되면, Snowpipe Streaming과 함께하는 Kafka 커넥터는 수동 상호작용이 필요합니다. Exactly-once 의미론은 여전히 보존됩니다.
- enable.streaming.client.optimization 속성이 false로 설정되어 있으면 Kafka 커넥터를 재시작해야 합니다. 커넥터를 재시작하면 새 주 배포를 대상으로 합니다.
- enable.streaming.client.optimization 속성이 true로 설정되어 있으면 커넥터가 실행 중인 호스트 JVM을 종료하고 다시 시작해야 합니다. 호스트 JVM을 재시작하면 새로 시작된 Kafka 커넥터가 새 주 배포를 대상으로 합니다.