Snowflake Connector for Kafka 동작 방식

Snowflake Connector for Kafka 동작 방식 (Working with the Snowflake Connector for Kafka)

이 문서는 Snowflake Connector for Kafka("Kafka 커넥터")가 어떻게 동작하는지, 그리고 주요 개념을 설명해요. 이 문서는 Kafka 커넥터 개요와 함께 읽으면 좋아요.

출처: Working with the Snowflake Connector for Kafka

본문

Kafka 커넥터의 동작 방식

Kafka 커넥터는 Kafka Connect 프레임워크에서 싱크 커넥터로 실행돼요. 커넥터는 다음 과정을 통해 Kafka 토픽의 데이터를 Snowflake 테이블로 로드해요:

  1. 토픽 구독: 커넥터는 구성 파일에 지정된 하나 이상의 Kafka 토픽을 구독해요.
  2. 메시지 버퍼링: 토픽 파티션에서 읽은 메시지는 메모리 버퍼에 누적돼요. 임계값(레코드 수, 시간, 크기)에 도달하면 메시지 일괄 처리가 시작돼요.
  3. 데이터 로드:
    • Snowpipe 모드: 메시지는 임시 파일로 스테이징된 후 Snowpipe가 로드해요.
    • Snowpipe Streaming 모드: 메시지는 Ingest SDK 버퍼를 통해 행 단위로 직접 Snowflake에 스트리밍돼요.
  4. 오프셋 커밋: 성공적으로 로드된 메시지의 오프셋이 커밋되어 재개 지점이 유지돼요.

수집 방식 선택

커넥터는 두 가지 수집 방식을 지원해요:

| 방식 | 설명 | 적합한 경우 | | Snowpipe | 파일 기반. 메시지를 임시 파일로 쓴 뒤 Snowpipe가 수집해요. | 일괄 처리, 큰 볼륨, 지연 허용이 큰 워크로드. | | Snowpipe Streaming | 행 기반. Ingest SDK를 통해 데이터를 거의 실시간으로 스트리밍해요. | 저지연, 실시간 분석이 필요한 워크로드. |

수집 방식은 snowflake.ingestion.method 구성 속성으로 지정해요 (SNOWPIPE 또는 SNOWPIPE_STREAMING).

커넥터가 만드는 Snowflake 객체

Kafka 커넥터는 각 토픽에 대해 다음 객체를 만들어요:

  • 내부 스테이지: 토픽 데이터 파일을 임시 저장하는 명명된 내부 스테이지 (Snowpipe 모드). 이름 형식: SNOWFLAKE_KAFKA_CONNECTOR_<connector_name>_STAGE_<table_name>.
  • 파이프 (Snowpipe 모드): 각 토픽 파티션에 대해 하나의 파이프. 이름 형식: SNOWFLAKE_KAFKA_CONNECTOR_<connector_name>_PIPE_<table_name>_<partition_number>.
  • 테이블: 각 토픽에 대해 하나의 테이블. 기본 스키마는 RECORD_CONTENT(VARIANT)와 RECORD_METADATA(VARIANT) 두 컬럼이에요.

이 객체들을 관리·정리하는 방법은 Kafka 커넥터 관리를 참고하세요.

테이블 스키마

기본으로 커넥터가 만드는 테이블:

  • RECORD_CONTENT: Kafka 메시지의 내용 (VARIANT).
  • RECORD_METADATA: 메시지에 대한 메타데이터 (VARIANT). 토픽, 파티션, 오프셋, createTime, 키, schema_id, 헤더 등을 포함해요.

RECORD_METADATA에 포함되는 필드:

| 필드 | 설명 | | topic | 레코드가 온 Kafka 토픽의 이름. | | partition | 토픽 내 파티션 번호. | | offset | 해당 파티션의 오프셋. | | CreateTime / LogAppendTime | Kafka 토픽의 메시지와 연결된 타임스탬프 (Unix epoch 밀리초). | | key | (KeyedMessage일 때) 메시지의 키. key.converter가 StringConverter일 때 저장돼요. | | schema_id | Avro 스키마 레지스트리 사용 시 스키마 ID. | | headers | 레코드와 연결된 헤더(키-값 쌍). |

스키마 감지 및 진화

커넥터는 snowflake.enable.schematization이 true(기본값)일 때 Kafka 메시지에서 필드를 감지하고 구조화된 스키마 컬럼을 만들 수 있어요. snowflake.schema.evolution을 true로 설정하면 테이블 스키마가 새 필드를 지원하도록 자동 진화해요.

장애 허용

Kafka와 Kafka 커넥터는 둘 다 장애 허용이 돼요. 메시지는 중복되거나 조용히 유실되지 않아요. Snowpipe 워크플로의 데이터 중복 제거 논리는 반복 데이터의 중복 복사본을 제거해요. Snowpipe가 레코드를 로드하는 중 오류를 만나면 레코드는 테이블 스테이지로 이동돼요.

Snowpipe Streaming 모드에서는 오류 처리를 위한 DLQ(dead-letter queue)를 지원해요.

오프셋 및 재개

커넥터는 처리한 오프셋을 추적해요. 커넥터가 중지되었다가 다시 시작되면 기록된 오프셋에서 재개해요. Snowpipe 모드에서는 내부 스테이지에 저장된 상태 정보를 사용해 재개 지점을 찾아요. 따라서 스테이지를 삭제하면 커넥터는 중단한 지점에서 재개할 수 없어요.

제한 사항

  • 커넥터 인스턴스는 서로 통신하지 않아요. 같은 토픽/파티션을 여러 인스턴스가 처리하면 중복 행이 삽입될 수 있어요. 각 토픽은 한 인스턴스만 처리해야 해요.
  • Kafka 토픽의 보존 기간(기본 7일)을 넘어 오프라인이면 만료된 레코드는 로드되지 않아요.
  • 메시지가 원래 발행된 순서대로 행이 삽입된다는 보장은 없어요.

더 알아보기 (Learn more)