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 커넥터")를 설치하고 구성하는 방법을 설명해요.
본문
전제 조건
Kafka 커넥터를 설치하려면 다음이 필요해요:
- Apache Kafka(버전 0.11 이상) 또는 Confluent 설치.
- Kafka Connect 클러스터.
- Java 8 이상.
- Snowflake 계정과 인증 방법(키 페어 인증 또는 OAuth).
설치
Kafka 커넥터를 설치하려면:
- 커넥터 JAR 파일을 다운로드하세요. Snowflake Kafka 커넥터 릴리스 페이지에서 다운로드할 수 있어요. Confluent Kafka를 사용한다면 Confluent Hub에서도 받을 수 있어요.
- JAR 파일을 Kafka 설치의
libs디렉토리에 복사하세요 (Confluent:<confluent_dir>/share/java/kafka-connect-snowflake/). - 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.SnowflakeJsonConvertercom.snowflake.kafka.connector.records.SnowflakeAvroConvertercom.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 커넥터를 업그레이드하려면:
- 새 버전의 JAR 파일을 다운로드하세요.
- 기존 JAR 파일을 새 파일로 교체하세요.
- Kafka Connect를 다시 시작하세요.