Apache Iceberg™ 테이블과 함께 Snowflake Connector for Kafka 사용

Apache Iceberg™ 테이블과 함께 Snowflake Connector for Kafka 사용 (Using the Snowflake Connector for Kafka with Apache Iceberg™ tables)

중요:

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

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

버전 3.0.0부터 Snowflake Connector for Kafka는 Snowflake 관리 Apache Iceberg™ 테이블로 데이터를 수집할 수 있어요.

출처: 문서

본문

요구사항 및 제한 사항

Iceberg 테이블 수집을 위해 Kafka 커넥터를 구성하기 전에 다음 요구사항과 제한 사항을 참고하세요:

  • Iceberg 테이블 수집에는 Kafka 커넥터 버전 3.0.0 이상이 필요해요.
  • Iceberg 테이블 수집은 Snowpipe Streaming을 사용하는 Kafka 커넥터에서 지원돼요. Snowpipe를 사용하는 Kafka 커넥터에서는 지원되지 않아요.
  • snowflake.streaming.enable.single.buffer가 false로 설정되면 Iceberg 테이블 수집이 지원되지 않아요.
  • 커넥터를 실행하기 전에 Iceberg 테이블을 만들어야 해요. 자세한 내용은 이 문서의 구성 및 설정을 참고하세요.

스키마 진화 제한 사항

Iceberg에 대한 스키마 진화는 AVRO나 Protobuf 같은 스키마가 있는 데이터 형식에 대해 완전히 지원돼요.

스키마가 없는 일반 JSON의 경우 커넥터는 다음 메시지 유형을 유효하지 않은 것으로 간주하고 DLQ(dead-letter queue)로 보내요:

  • 해당 값이 null 또는 []인 새 컬럼이 있는 메시지.
  • 해당 값이 null 또는 []인 구조화 객체의 새 필드가 있는 메시지.

커넥터가 이러한 메시지 유형을 수집할 수 있도록 테이블 스키마를 수동으로 변경하려면 ALTER TABLE 문을 사용하세요.

구성 및 설정

Iceberg 테이블 수집용으로 Kafka 커넥터를 구성하려면 다음 섹션에 명시된 몇 가지 차이점을 제외하고 Snowpipe Streaming 기반 커넥터의 일반 설정 단계를 따르세요.

Iceberg 테이블은 다음 스토리지 옵션 중 하나를 사용할 수 있어요:

  • Snowflake 스토리지 (영구 Iceberg 테이블만): Snowflake가 Iceberg 테이블 파일을 만들고 관리해줘서 외부 볼륨을 만들거나 커넥터에 접근 권한을 부여할 필요가 없어요. Snowflake 스토리지를 사용하는 임시(transient) Iceberg 테이블은 지원되지 않아요. 자세한 내용은 클래식 아키텍처와 함께 Snowpipe Streaming의 제한 사항 및 고려사항을 참고하세요.
  • 사용자가 관리하는 외부 클라우드 스토리지 — 외부 볼륨을 통해 접근. 커넥터 역할에 외부 볼륨에 대한 USAGE 권한을 부여해야 해요.

외부 볼륨에 사용 권한 부여

이 단계는 Iceberg 테이블이 사용자가 관리하는 외부 볼륨을 사용할 때만 적용돼요. 테이블이 Snowflake 스토리지를 사용하면 이 단계를 건너뛰세요.

예를 들어 Iceberg 테이블이 kafka_external_volume 외부 볼륨을 사용하고 커넥터가 kafka_connector_role 역할을 사용한다면 다음 문을 실행하세요:

USE ROLE ACCOUNTADMIN;
GRANT USAGE ON EXTERNAL VOLUME kafka_external_volume TO ROLE kafka_connector_role;

수집용 Iceberg 테이블 생성

커넥터를 실행하기 전에 Iceberg 테이블을 만들어야 해요. 초기 테이블 스키마는 커넥터의 snowflake.enable.schematization 설정에 따라 달라져요.

스키마화(schematization)를 활성화하면 record_metadata라는 컬럼으로 테이블을 만들 수 있어요. Snowflake 스토리지를 사용하려면 EXTERNAL_VOLUME = 'SNOWFLAKE_MANAGED'를 설정하고 BASE_LOCATION을 생략하세요:

CREATE OR REPLACE ICEBERG TABLE my_iceberg_table (
    record_metadata OBJECT()
  )
  EXTERNAL_VOLUME = 'SNOWFLAKE_MANAGED'
  CATALOG = 'SNOWFLAKE';

자체 외부 볼륨을 사용하려면 EXTERNAL_VOLUME을 볼륨 이름으로 설정하고 BASE_LOCATION을 제공하세요:

CREATE OR REPLACE ICEBERG TABLE my_iceberg_table (
    record_metadata OBJECT()
  )
  EXTERNAL_VOLUME = 'my_volume'
  CATALOG = 'SNOWFLAKE'
  BASE_LOCATION = 'my_location/my_iceberg_table';

커넥터는 메시지 필드에 대한 컬럼을 자동으로 만들고 record_metadata 컬럼 스키마를 변경해요.

스키마화를 활성화하지 않으면 실제 Kafka 메시지 콘텐츠와 일치하는 타입의 record_content라는 컬럼으로 테이블을 만들 수 있어요. 커넥터는 record_metadata 컬럼을 자동으로 만들어요.

Iceberg 테이블을 만들 때 Iceberg 데이터 타입이나 호환 가능한 Snowflake 타입을 사용할 수 있어요. 반정형 VARIANT 타입은 지원되지 않아요. 대신 구조화된 OBJECT 또는 MAP을 사용하세요.

예를 들어 다음 메시지를 고려해 보세요:

{
    "id": 1,
    "name": "Steve",
    "body_temperature": 36.6,
    "approved_coffee_types": ["Espresso", "Doppio", "Ristretto", "Lungo"],
    "animals_possessed":
    {
        "dogs": true,
        "cats": false
    },
    "date_added": "2024-10-15"
}

예시 메시지용 Iceberg 테이블을 만들려면 다음 문을 사용하세요:

CREATE OR REPLACE ICEBERG TABLE my_iceberg_table (
    record_content OBJECT(
        id INT,
        body_temperature FLOAT,
        name STRING,
        approved_coffee_types ARRAY(STRING),
        animals_possessed OBJECT(dogs BOOLEAN, cats BOOLEAN),
        date_added DATE
    )
  )
  EXTERNAL_VOLUME = 'my_volume'
  CATALOG = 'SNOWFLAKE'
  BASE_LOCATION = 'my_location/my_iceberg_table';

참고: dogs나 cats 같은 중첩 구조 내부의 필드 이름은 대소문자를 구분해요.

구성 속성

snowflake.streaming.iceberg.enabled — 커넥터가 Iceberg 테이블에 데이터를 수집하는지 여부를 지정해요. 이 속성이 실제 테이블 유형과 일치하지 않으면 커넥터가 실패해요.

값:

  • true
  • false

기본값: false

더 알아보기 (Learn more)