Kafka 커넥터 설치 및 구성

Kafka 커넥터 설치 및 구성 (Installing and configuring the Kafka connector)

중요:

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

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

이 문서는 Snowflake Connector for Kafka("Kafka 커넥터")를 설치하고 구성하는 방법을 설명해요.

출처: Installing and configuring the Kafka connector

본문

전제 조건

Kafka 커넥터를 설치하려면 다음이 필요해요:

설치

Kafka 커넥터를 설치하려면:

  1. 커넥터 JAR 파일을 다운로드하세요. Snowflake Kafka 커넥터 릴리스 페이지에서 다운로드할 수 있어요. Confluent Kafka를 사용한다면 Confluent Hub에서도 받을 수 있어요.
  2. JAR 파일을 Kafka 설치의 libs 디렉토리에 복사하세요 (Confluent: <confluent_dir>/share/java/kafka-connect-snowflake/).
  3. Kafka Connect를 다시 시작하세요.

Kafka 커넥터 구성

커넥터를 구성하려면 커넥터의 구성 파일(예: <kafka_dir>/config/connect-distributed-snowflake.properties)에 속성을 설정해야 해요. 다음은 최소 필수 속성의 예시예요:

name=XYZCompanySensorData
connector.class=com.snowflake.kafka.connector.SnowflakeSinkConnector
tasks.max=8
topics=topic1,topic2
snowflake.url.name=myorganization-myaccount.snowflakecomputing.com:443
snowflake.user.name=jane.smith
snowflake.private.key=xyz123
snowflake.private.key.passphrase=jkladu098jfd089adsq4r
snowflake.database.name=mydb
snowflake.schema.name=myschema
buffer.count.records=10000
buffer.flush.time=60
buffer.size.bytes=5000000
snowflake.topic2table.map=topic1:table1,topic2:table2

필수 속성

| 속성 | 데이터 타입 | 필수 | 설명 | | name | String | 예 | 커넥터의 이름. | | connector.class | String | 예 | 커넥터 클래스: com.snowflake.kafka.connector.SnowflakeSinkConnector. | | tasks.max | Integer | 예 | 커넥터에서 사용할 태스크의 최대 수. 토픽 파티션 수를 기준으로 권장값은 1에서 4 사이예요. | | topics | String | 예 | Kafka 커넥터가 데이터를 읽을 토픽 목록(쉼표로 구분). | | snowflake.url.name | String | 예 | Snowflake 계정 이름과 계정 로케이터 형식 (<organization_name>-<account_name>.snowflakecomputing.com:443). | | snowflake.user.name | String | 예 | Snowflake 사용자 이름. | | snowflake.private.key | String | 조건부 | 키 페어 인증을 사용하는 경우 프라이빗 키 (PEM 형식, Base64 인코딩된 DER 형식, 또는 PKCS#8 형식). | | snowflake.private.key.passphrase | String | 조건부 | 프라이빗 키가 암호화된 경우 해당 암호(passphrase). | | snowflake.database.name | String | 예 | 데이터를 로드할 Snowflake 데이터베이스. | | snowflake.schema.name | String | 예 | 데이터를 로드할 Snowflake 스키마. |

참고: OAuth를 사용해 인증하려면 snowflake.user.name 대신 snowflake.oauth.client.id와 snowflake.oauth.client.secret 속성을 설정할 수 있어요. 자세한 내용은 Snowflake Kafka 커넥터로 OAuth 구성을 참고하세요.

선택적 속성

| 속성 | 데이터 타입 | 기본값 | 설명 | | snowflake.ingestion.method | String | SNOWPIPE | Snowpipe(SNOWPIPE) 또는 Snowpipe Streaming(SNOWPIPE_STREAMING) 중 데이터 로딩 방법을 지정해요. | | buffer.count.records | Integer | 10000 | 버퍼가 플러시되기 전에 누적할 레코드 수. | | buffer.flush.time | Integer | 60 | 버퍼가 플러시되기 전의 최대 시간(초). | | buffer.size.bytes | Long | 5000000 | 버퍼가 플러시되기 전의 최대 크기(바이트). | | snowflake.topic2table.map | String | | 토픽을 테이블에 매핑. 형식: topic1:table1,topic2:table2. 지정하지 않으면 커넥터가 각 토픽에 대해 같은 이름의 테이블을 생성해요. | | snowflake.role.name | String | | 커넥터가 사용할 Snowflake 역할. | | snowflake.schema.evolution | Boolean | false | 스키마 진화 활성화 여부. | | snowflake.enable.schematization | Boolean | true | 메시지를 구조화된 스키마로 변환할지 여부. | | key.converter | String | | Kafka 메시지 키에 대한 변환기 클래스. | | value.converter | String | | Kafka 메시지 값에 대한 변환기 클래스. | | jmx | Boolean | true | JMX 활성화 여부. |

Snowpipe Streaming 관련 속성

Snowpipe Streaming을 사용하면 다음 속성을 추가로 구성할 수 있어요:

| 속성 | 데이터 타입 | 기본값 | 설명 | | snowflake.streaming.client.auto.buffer.flush.time | Integer | 5 | 버퍼가 자동으로 플러시되기 전 최대 시간(초). | | snowflake.streaming.enable.single.buffer | Boolean | true | 단일 버퍼 활성화 여부. Iceberg 테이블 수집을 사용하려면 true여야 해요. | | snowflake.streaming.max.client.errors | Integer | 10 | 클라이언트가 버퍼를 계속 사용하기 전 허용되는 최대 오류 수. | | snowflake.streaming.min.buffer.time | Integer | 5 | 버퍼가 플러시되기 전 최소 대기 시간(초). | | snowflake.streaming.partitioning.enabled | Boolean | true | 파티셔닝 활성화 여부. |

테이블 만들기

기본적으로 Kafka 커넥터는 각 토픽에 대해 RECORD_CONTENT(VARIANT)와 RECORD_METADATA(VARIANT) 두 컬럼을 가진 테이블을 자동으로 생성해요. 테이블을 미리 만들고 싶다면 다음 DDL로 만들 수 있어요:

CREATE OR REPLACE TABLE mytable (
  RECORD_CONTENT VARIANT,
  RECORD_METADATA VARIANT
);

KEY 및 VALUE 형식

Kafka 커넥터는 메시지의 key와 value를 변환하기 위해 변환기(converter)를 사용해요. 기본적으로 Snowflake 변환기가 사용돼요:

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

Avro와 Protobuf는 스키마 레지스트리와 함께 사용할 수 있어요. 자세한 내용은 protobuf 데이터 로드와 Avro 데이터 로드를 참고하세요.

스키마 감지 및 진화

snowflake.enable.schematization이 true(기본값)로 설정되면 커넥터는 각 Kafka 메시지에서 필드를 감지해 테이블에 구조화된(structured) 컬럼을 만드려고 시도해요. snowflake.schema.evolution이 true로 설정되면 테이블 스키마는 새 필드를 지원하도록 자동 진화해요. 감지할 수 없는 메시지나 잘못된 스키마의 메시지는 표준 구조의 테이블(두 VARIANT 컬럼)로 로드될 수 있어요.

스키마 감지가 비활성화된 경우

snowflake.enable.schematization을 false로 설정하면 모든 메시지가 RECORD_CONTENT(VARIANT)와 RECORD_METADATA(VARIANT) 컬럼으로 로드돼요.

인증

Kafka 커넥터는 다음 인증 방법을 지원해요:

  • 키 페어 인증: snowflake.private.key 속성으로 프라이빗 키를 지정. 프라이빗 키는 암호화되지 않은 PEM, 암호화된 PEM, Base64 인코딩된 DER, 또는 PKCS#8 형식일 수 있어요. 자세한 내용은 키 페어 인증을 참고하세요.
  • OAuth: snowflake.oauth.client.id와 snowflake.oauth.client.secret 속성을 설정해 OAuth로 인증. 자세한 내용은 Kafka 커넥터 OAuth 구성을 참고하세요.

역할 및 권한

Kafka 커넥터가 사용하는 Snowflake 사용자에게 다음 권한을 부여해야 해요 (예: ACCOUNTADMIN 역할 또는 커스텀 역할):

GRANT USAGE ON DATABASE mydb TO ROLE kafka_connector_role;
GRANT USAGE ON SCHEMA mydb.myschema TO ROLE kafka_connector_role;
GRANT CREATE STAGE ON SCHEMA mydb.myschema TO ROLE kafka_connector_role;
GRANT CREATE PIPE ON SCHEMA mydb.myschema TO ROLE kafka_connector_role;
GRANT CREATE TABLE ON SCHEMA mydb.myschema TO ROLE kafka_connector_role;
GRANT CREATE FILE FORMAT ON SCHEMA mydb.myschema TO ROLE kafka_connector_role;
GRANT INSERT ON TABLE mydb.myschema.mytable TO ROLE kafka_connector_role;
GRANT OWNERSHIP ON TABLE mydb.myschema.mytable TO ROLE kafka_connector_role;

참고: Snowpipe를 사용하는 경우 스테이지에 대한 OWNERSHIP 권한과 파이프에 대한 OWNERSHIP 권한도 필요해요.

Snowpipe 사용 시 자동 파이프 생성

snowflake.ingestion.method=SNOWPIPE로 설정하면 커넥터가 각 토픽 파티션에 대해 자동으로 파이프를 생성해요. 파이프 이름 형식은 SNOWFLAKE_KAFKA_CONNECTOR_<connector_name>_PIPE_<table_name>_<partition_number>이에요. 커넥터가 만든 객체를 관리하려면 Kafka 커넥터 관리를 참고하세요.

로드된 데이터 확인

데이터가 로드되었는지 확인하려면 테이블을 쿼리해 보세요:

SELECT * FROM mydb.myschema.mytable LIMIT 10;

Kafka 커넥터 업그레이드

Kafka 커넥터를 업그레이드하려면:

  1. 새 버전의 JAR 파일을 다운로드하세요.
  2. 기존 JAR 파일을 새 파일로 교체하세요.
  3. Kafka Connect를 다시 시작하세요.

더 알아보기 (Learn more)