Kafka 커넥터 개요

Kafka 커넥터 개요 (Overview of the Kafka connector)

중요:

  • 사전 공지: 클래식 Kafka 커넥터(v3 이하)는 현재 완전히 지원되지만, 향후 폐기될 예정이에요.
  • 조치: 즉시 변경할 필요는 없어요. 현재 워크로드는 안전하며 계속 완전히 지원돼요.
  • 일정: Snowflake는 2026년 중반에 공식 폐기 공지를 발표할 계획이에요. 공지 후 수명 종료까지 18개월의 마이그레이션 기간이 시작돼요.
  • 권장사항: 모든 새 구현에는 Snowflake Connector for Kafka (v4)를 사용하세요.

마이그레이션 지침은 v3에서 v4로 마이그레이션을 참고하세요.

이 문서는 Apache Kafka와 Snowflake Connector for Kafka에 대한 개요를 제공해요.

참고: Kafka 커넥터는 커넥터 약관이 적용돼요.

출처: Overview of the Kafka connector

본문

Apache Kafka 소개

Apache Kafka 소프트웨어는 메시지 큐나 엔터프라이즈 메시징 시스템과 유사하게 레코드 스트림을 쓰고 읽는 발행·구독 모델을 사용해요. Kafka는 프로세스가 메시지를 비동기적으로 읽고 쓸 수 있게 해줘요. 구독자는 발행자에 직접 연결될 필요가 없어요. 발행자는 구독자가 나중에 받도록 Kafka에 메시지를 큐잉할 수 있어요.

애플리케이션은 **토픽(topic)**에 메시지를 발행하고, 애플리케이션은 토픽을 구독해 그 메시지를 받아요. Kafka는 메시지를 전송할 뿐만 아니라 처리할 수도 있지만, 이는 이 문서의 범위 밖이에요. 토픽은 확장성을 높이기 위해 **파티션(partition)**으로 나눌 수 있어요.

Kafka Connect는 Kafka를 데이터베이스를 포함한 외부 시스템과 연결하는 프레임워크예요. Kafka Connect 클러스터는 Kafka 클러스터와 분리된 별도의 클러스터예요. Kafka Connect 클러스터는 커넥터(외부 시스템과의 읽기/쓰기 지원 컴포넌트)의 실행과 확장을 지원해요.

Kafka 커넥터는 Kafka Connect 클러스터에서 실행되어 Kafka 토픽에서 데이터를 읽고 해당 데이터를 Snowflake 테이블에 쓰도록 설계됐어요.

Snowflake는 두 가지 버전의 커넥터를 제공해요:

Snowflake의 관점에서 Kafka 토픽은 Snowflake 테이블에 삽입될 행 스트림을 생성해요. 일반적으로 각 Kafka 메시지는 하나의 행을 포함해요.

Kafka는 많은 메시지 발행/구독 플랫폼처럼 발행자와 구독자 간의 다대다 관계를 허용해요. 단일 애플리케이션이 여러 토픽에 발행할 수 있고, 단일 애플리케이션이 여러 토픽을 구독할 수 있어요. Snowflake에서 일반적인 패턴은 하나의 토픽이 하나의 Snowflake 테이블에 메시지(행)를 공급하는 것이에요.

Kafka 커넥터의 현재 버전은 Snowflake로의 데이터 로딩에 제한돼 있어요. Kafka 커넥터는 두 가지 데이터 로딩 방법을 지원해요:

  • Snowpipe
  • Snowpipe Streaming

자세한 내용은 Snowflake로 데이터 로드와 Snowpipe Streaming으로 Snowflake Connector for Kafka 사용을 참고하세요.

Kafka 토픽의 대상 테이블

Kafka 토픽은 Kafka 구성에서 기존 Snowflake 테이블에 매핑할 수 있어요. 토픽이 매핑되지 않으면 Kafka 커넥터가 토픽 이름을 사용해 각 토픽에 대한 새 테이블을 만들어요.

커넥터는 다음 규칙을 사용해 토픽 이름을 유효한 Snowflake 테이블 이름으로 변환해요:

  • 소문자 토픽 이름은 대문자 테이블 이름으로 변환돼요.
  • 토픽 이름의 첫 문자가 문자(a-z 또는 A-Z)나 밑줄 문자(_)가 아니면 커넥터는 테이블 이름 앞에 밑줄을 붙여요.
  • 토픽 이름 안의 어떤 문자가 Snowflake 테이블 이름에 유효하지 않은 문자라면 그 문자는 밑줄 문자로 교체돼요. 테이블 이름에 유효한 문자에 대한 자세한 내용은 식별자 요구사항을 참고하세요.

Kafka 커넥터가 Kafka 토픽용으로 만든 테이블의 이름을 조정해야 한다면 같은 스키마에 두 테이블의 이름이 동일해질 수 있어요. 예를 들어 numbers+x와 numbers-x 토픽에서 데이터를 읽으면 이 토픽들을 위해 만들어진 테이블이 둘 다 NUMBERS_X가 돼요. 우발적인 테이블 이름 중복을 피하기 위해 커넥터는 테이블 이름에 접미사를 붙여요. 접미사는 밑줄 뒤에 생성된 해시 코드가 붙는 형태예요.

팁: 가능하다면 Snowflake 식별자 이름 규칙을 따르는 토픽 이름을 선택하는 것을 권장해요.

Kafka 토픽용 테이블의 스키마

Kafka 커넥터가 로드하는 테이블의 스키마는 테이블 유형과 커넥터 구성 방법에 따라 달라져요:

  • 대부분의 테이블은 이 섹션에 설명된 기본 스키마를 사용해요.
  • 스키마 감지 및 진화를 사용하면 스키마에 사용자 정의 스키마와 일치하는 컬럼이 포함돼요.
  • Iceberg 테이블로 수집하면 스키마에 동일한 기본 컬럼(record_content와 record_metadata)이 포함돼요. 다만 VARIANT 대신 구조화된 타입 컬럼이에요.

기본적으로 Snowpipe나 Snowpipe Streaming으로 Kafka 커넥터가 로드하는 모든 Snowflake 테이블에는 두 개의 VARIANT 컬럼으로 구성된 스키마가 있어요:

  • RECORD_CONTENT: Kafka 메시지를 포함해요.
  • RECORD_METADATA: 메시지에 대한 메타데이터(예를 들어 메시지를 읽은 토픽)를 포함해요.

Snowflake가 테이블을 만들면 테이블에는 이 두 컬럼만 포함돼요. 사용자가 Kafka 커넥터가 행을 추가할 테이블을 만들면 테이블에 이 두 컬럼보다 더 많은 컬럼이 포함될 수 있어요 (추가 컬럼은 커넥터의 데이터가 해당 컬럼에 대한 값을 포함하지 않으므로 NULL을 허용해야 해요).

RECORD_CONTENT 컬럼은 Kafka 메시지를 포함해요.

Kafka 메시지는 전송되는 정보에 따라 달라지는 내부 구조를 가져요. 예를 들어 IoT(사물 인터넷) 날씨 센서의 메시지에는 데이터가 기록된 타임스탬프, 센서 위치, 온도, 습도 등을 포함할 수 있어요. 재고 시스템의 메시지에는 제품 ID와 판매된 품목 수, 판매·배송 시각을 나타내는 타임스탬프를 포함할 수 있어요.

일반적으로 특정 토픽의 각 메시지는 같은 기본 구조를 가져요. 다른 토픽은 일반적으로 다른 구조를 사용해요.

각 Kafka 메시지는 JSON 형식 또는 Avro 형식으로 Snowflake에 전달돼요. Kafka 커넥터는 이 형식화된 정보를 VARIANT 타입의 단일 컬럼에 저장해요. 데이터는 파싱되지 않고 Snowflake 테이블에서 여러 컬럼으로 분리되지 않아요.

RECORD_METADATA 컬럼은 기본적으로 다음 정보를 포함해요:

| 필드 | Java 데이터 타입 | SQL 데이터 타입 | 필수 | 설명 | | topic | String | VARCHAR | 예 | 레코드가 온 Kafka 토픽의 이름. | | partition | String | VARCHAR | 예 | 토픽 내 파티션 번호. (이것은 Kafka 파티션이지 Snowflake 마이크로 파티션이 아니에요.) | | offset | long | INTEGER | 예 | 해당 파티션의 오프셋. | | CreateTime / LogAppendTime | long | BIGINT | 아니요 | Kafka 토픽의 메시지와 연결된 타임스탬프. 값은 1970년 1월 1일 자정(UTC) 이후 밀리초 단위예요. 자세한 내용: https://kafka.apache.org/0100/javadoc/org/apache/kafka/clients/producer/ProducerRecord.html | | SnowflakeConnectorPushTime | long | BIGINT | 아니요 | Snowpipe Streaming 사용 시에만 사용 가능. 레코드가 Ingest SDK 버퍼로 푸시된 시각의 타임스탬프. 값은 1970년 1월 1일 자정(UTC) 이후 밀리초 수예요. 자세한 내용: 수집 지연 시간 추정 | | key | String | VARCHAR | 아니요 | 메시지가 Kafka KeyedMessage이면 이 메시지의 키. 커넥터가 키를 RECORD_METADATA에 저장하려면 Kafka 구성 속성의 key.converter 파라미터가 "org.apache.kafka.connect.storage.StringConverter"로 설정되어야 해요. 그렇지 않으면 커넥터가 키를 무시해요. | | schema_id | int | INTEGER | 아니요 | 스키마 레지스트리로 Avro를 사용해 스키마를 지정하면 이 레지스트리의 스키마 ID. | | headers | Object | OBJECT | 아니요 | 헤더는 레코드와 연결된 사용자 정의 키-값 쌍이에요. 각 레코드는 0, 1 또는 여러 헤더를 가질 수 있어요. |

RECORD_METADATA 컬럼에 기록되는 메타데이터의 양은 선택적 Kafka 구성 속성을 사용해 구성할 수 있어요. 자세한 내용은 Kafka 커넥터 설치 및 구성을 참고하세요.

필드 이름과 값은 대소문자를 구분해요.

JSON 구문으로 표현하면 샘플 메시지는 다음과 유사할 수 있어요:

{
    "meta":
    {
        "offset": 1,
        "topic": "PressureOverloadWarning",
        "partition": 12,
        "key": "key name",
        "schema_id": 123,
        "CreateTime": 1234567890,
        "headers":
        {
            "name1": "value1",
            "name2": "value2"
        }
    },
    "content":
    {
        "ID": 62,
        "PSI": 451,
        "etc": "..."
    }
}

VARIANT 컬럼 쿼리 구문을 사용해 Snowflake 테이블을 직접 쿼리할 수 있어요.

RECORD_METADATA의 토픽을 기반으로 데이터를 추출하는 간단한 예시:

select
       record_metadata:CreateTime,
       record_content:ID
    from table1
    where record_metadata:topic = 'PressureOverloadWarning';

출력은 다음과 유사해요:

+------------+-----+
| CREATETIME | ID  |
+------------+-----+
| 1234567890 | 62  |
+------------+-----+

또는 이 테이블에서 데이터를 추출해 개별 컬럼으로 펼치고(flatten) 보통 쿼리하기 더 쉬운 다른 테이블에 저장할 수 있어요.

Kafka 커넥터 워크플로

Kafka 커넥터는 Kafka 토픽을 구독하고 Snowflake 객체를 생성하기 위해 다음 프로세스를 완료해요:

  • Kafka 커넥터는 Kafka 구성 파일이나 명령줄(또는 Confluent Control Center; Confluent만)을 통해 제공된 구성 정보를 기반으로 하나 이상의 Kafka 토픽을 구독해요.
  • 커넥터는 각 토픽에 대해 다음 객체를 생성해요:
    • 각 토픽의 데이터 파일을 임시로 저장하는 내부 스테이지 하나.
    • 각 토픽 파티션의 데이터 파일을 수집하는 파이프 하나.
    • 각 토픽에 대한 테이블 하나. 각 토픽에 지정된 테이블이 존재하지 않으면 커넥터가 생성하고, 그렇지 않으면 기존 테이블에 RECORD_CONTENT와 RECORD_METADATA 컬럼을 만들고 다른 컬럼이 nullable인지 확인해요 (null이 아니면 오류를 생성해요).

다음 다이어그램은 Kafka 커넥터를 사용한 Kafka 수집 흐름을 보여줘요:

  • 하나 이상의 애플리케이션이 JSON 또는 Avro 레코드를 Kafka 클러스터에 발행해요. 레코드는 하나 이상의 토픽 파티션으로 분리돼요.
  • Kafka 커넥터는 Kafka 토픽의 메시지를 버퍼링해요. 임계값(시간, 메모리, 또는 메시지 수)에 도달하면 커넥터는 메시지를 내부 스테이지의 임시 파일에 써요. 커넥터는 Snowpipe를 트리거해 임시 파일을 수집해요. Snowpipe는 데이터 파일에 대한 포인터를 큐에 복사해요.
  • Snowflake 제공 가상 웨어하우스가 Kafka 토픽 파티션용으로 만들어진 파이프를 통해 스테이징된 파일에서 대상 테이블(즉 구성 파일에서 토픽에 대해 지정된 테이블)로 데이터를 로드해요.
  • (표시되지 않음) 커넥터는 Snowpipe를 모니터링하고 파일 데이터가 테이블에 로드되었음을 확인한 후 내부 스테이지의 각 파일을 삭제해요. 실패로 데이터가 로드되지 않으면 커넥터는 파일을 테이블 스테이지로 옮기고 오류 메시지를 생성해요.
  • 커넥터는 2~4단계를 반복해요.

주의: Snowflake는 insertReport API를 1시간 동안 폴링해요. 이 시간 내에 수집된 파일의 상태가 성공하지 않으면 수집 중인 파일이 테이블 스테이지로 이동돼요. 이 파일들이 테이블 스테이지에서 사용 가능해지려면 최소 1시간이 걸릴 수 있어요. 파일은 수집 상태를 이전 1시간 내에 찾을 수 없을 때만 테이블 스테이지로 이동돼요.

장애 허용 (Fault tolerance)

Kafka와 Kafka 커넥터는 둘 다 장애 허용이 돼요. 메시지는 중복되거나 조용히 유실되지 않아요.

데이터 로딩 체인의 Snowpipe 워크플로의 데이터 중복 제거 로직은 드문 경우를 제외하고 반복 데이터의 중복 복사본을 제거해요. Snowpipe가 레코드를 로드하는 동안 오류가 감지되면(예: 레코드가 잘 구성된 JSON이나 Avro가 아님) 레코드는 로드되지 않고 테이블 스테이지로 이동돼요.

Snowpipe Streaming을 사용하는 Kafka 커넥터는 오류 처리를 위한 DLQ(dead-letter queue)를 지원해요. 자세한 내용은 Snowpipe Streaming을 사용하는 Kafka 커넥터의 오류 처리 및 DLQ 속성을 참고하세요.

커넥터의 장애 허용 제한

Kafka 토픽은 스토리지 공간 또는 보존 시간에 제한을 두도록 구성할 수 있어요.

  • 기본 보존 시간은 7일이에요. 시스템이 보존 시간보다 오래 오프라인 상태면 만료된 레코드는 로드되지 않아요. 마찬가지로 Kafka의 스토리지 공간 제한을 초과하면 일부 메시지가 전달되지 않아요.
  • Kafka 토픽의 메시지가 삭제되거나 업데이트되면 이러한 변경 사항이 Snowflake 테이블에 반영되지 않을 수 있어요.

주의: Kafka 커넥터의 인스턴스는 서로 통신하지 않아요. 같은 토픽이나 파티션에서 커넥터의 여러 인스턴스를 시작하면 같은 행의 여러 복사본이 테이블에 삽입될 수 있어요. 권장되지 않아요. 각 토픽은 커넥터의 인스턴스 하나만 처리해야 해요.

이론적으로는 Kafka에서 Snowflake가 수집할 수 있는 것보다 빠르게 메시지가 흐를 수 있어요. 그러나 실제로는 이 경우가 드물어요. 발생한다면 Kafka Connect 클러스터의 성능 튜닝으로 문제를 해결해야 해요. 예를 들어:

  • Connect 클러스터의 노드 수 튜닝.
  • 커넥터에 할당된 태스크 수 튜닝.
  • 커넥터와 Snowflake 배포 간 네트워크 대역폭의 영향 이해.

중요: 원래 발행된 순서대로 행이 삽입된다는 보장은 없어요.

지원되는 플랫폼

Kafka 커넥터는 모든 Kafka Connect 클러스터에서 실행될 수 있고, 지원되는 클라우드 플랫폼의 Snowflake 계정으로 데이터를 보낼 수 있어요.

Protobuf 데이터 지원

Kafka 커넥터 1.5.0(이상)은 protobuf 변환기를 통해 프로토콜 버퍼(protobuf)를 지원해요. 자세한 내용은 Snowflake Connector for Kafka로 protobuf 데이터 로드를 참고하세요.

청구 정보

Kafka 커넥터 사용에 대한 직접적인 요금은 없어요. 그러나 간접 비용이 있어요:

  • Snowpipe는 커넥터가 Kafka에서 읽은 데이터를 로드하는 데 사용되며, Snowpipe 처리 시간은 계정에 청구돼요.
  • 데이터 스토리지는 계정에 청구돼요.

Kafka 커넥터 제한 사항

SMT(Single Message Transformation)는 메시지가 Kafka Connect를 통과할 때 적용돼요. Kafka 구성 속성을 구성할 때 key.converter나 value.converter 중 하나를 다음 값 중 하나로 설정하면 해당 키나 값에 SMT가 지원되지 않아요:

  • com.snowflake.kafka.connector.records.SnowflakeJsonConverter
  • com.snowflake.kafka.connector.records.SnowflakeAvroConverter
  • com.snowflake.kafka.connector.records.SnowflakeAvroConverterWithoutSchemaRegistry

key.converter와 value.converter 모두 설정되지 않으면 대부분의 SMT가 지원되며, 현재 regex.router는 예외예요.

Snowflake 변환기가 SMT를 지원하지 않지만, Kafka 커넥터 버전 1.4.3(이상)은 다음과 같은 많은 커뮤니티 기반 변환기를 지원해요:

  • io.confluent.connect.avro.AvroConverter
  • org.apache.kafka.connect.json.JsonConverter

SMT에 대한 자세한 내용은 https://docs.confluent.io/current/connect/transforms/index.html를 참고하세요.

더 알아보기 (Learn more)