Snowflake Connector for Kafka로 protobuf 데이터 로드

Snowflake Connector for Kafka로 protobuf 데이터 로드 (Loading protobuf data using the Snowflake Connector for Kafka)

중요:

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

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

이 문서는 Snowflake Connector for Kafka("Kafka 커넥터")에 프로토콜 버퍼(protobuf) 지원을 설치·구성하는 방법을 설명해요. protobuf 지원에는 Kafka 커넥터 1.5.0(이상)이 필요해요.

Kafka 커넥터는 다음 버전의 protobuf 변환기를 지원해요:

  • Confluent 버전: Confluent 패키지 버전의 Kafka에서만 지원돼요.
  • 커뮤니티 버전: OSS(오픈소스 소프트웨어) Apache Kafka 패키지에서 지원돼요. Confluent 패키지 버전의 Kafka에서도 지원되지만, 사용 편의를 위해 Confluent 버전을 사용하는 것을 권장해요.

이 protobuf 변환기 중 하나만 설치하세요.

출처: Loading protobuf data using the Snowflake Connector for Kafka

본문

전제 조건: Snowflake Connector for Kafka 설치

Kafka 커넥터 설치 및 구성의 지침을 사용해 Kafka 커넥터를 설치하세요.

Confluent 버전의 protobuf 변환기 구성

참고: Protobuf 변환기의 Confluent 버전은 Confluent 버전 5.5.0(이상)에서 사용할 수 있어요.

  • 텍스트 편집기에서 Kafka 구성 파일(예: <kafka_dir>/config/connect-distributed.properties)을 여세요.
  • 파일에서 변환기 속성을 구성하세요. 일반적인 Kafka 커넥터 속성에 대한 내용은 Kafka 구성 속성을 참고하세요.
{
 "name":"XYZCompanySensorData",
"config":{
  ..
  "key.converter":"io.confluent.connect.protobuf.ProtobufConverter",
  "key.converter.schema.registry.url":"CONFLUENT_SCHEMA_REGISTRY",
  "value.converter":"io.confluent.connect.protobuf.ProtobufConverter",
  "value.converter.schema.registry.url":"http://localhost:8081"
}
 }

예를 들어:

{
  "name":"XYZCompanySensorData",
  "config":{
 "connector.class":"com.snowflake.kafka.connector.SnowflakeSinkConnector",
 "tasks.max":"8",
 "topics":"topic1,topic2",
 "snowflake.topic2table.map": "topic1:table1,topic2:table2",
 "buffer.count.records":"10000",
 "buffer.flush.time":"60",
 "buffer.size.bytes":"5000000",
 "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",
 "key.converter":"io.confluent.connect.protobuf.ProtobufConverter",
 "key.converter.schema.registry.url":"CONFLUENT_SCHEMA_REGISTRY",
 "value.converter":"io.confluent.connect.protobuf.ProtobufConverter",
 "value.converter.schema.registry.url":"http://localhost:8081"
  }
}
  • 파일을 저장하세요.

Confluent 콘솔 protobuf 프로듀서, 소스 protobuf 프로듀서, 또는 Python 프로듀서를 사용해 Kafka에서 protobuf 데이터를 생성(produce)하세요. GitHub의 예제 Python 코드는 Kafka에서 protobuf 데이터를 생성하는 방법을 보여줘요.

커뮤니티 버전의 protobuf 변환기 구성

이 섹션은 커뮤니티 버전의 protobuf 변환기를 설치·구성하는 방법을 설명해요.

1단계: 커뮤니티 protobuf 변환기 설치

  • 터미널 창에서 protobuf 변환기의 GitHub 저장소 클론을 저장할 디렉토리로 변경하세요.
  • 다음 명령을 실행해 GitHub 저장소를 클론하세요:
git clone https://github.com/blueapron/kafka-connect-protobuf-converter
  • 다음 명령을 실행해 Apache Maven으로 변환기의 3.1.0 버전을 빌드하세요. Kafka 커넥터는 변환기의 2.3.0, 3.0.0, 3.1.0 버전을 지원한다는 점에 유의하세요. (참고: Maven이 로컬 머신에 이미 설치되어 있어야 해요.)
cd kafka-connect-protobuf-converter

git checkout tags/v3.1.0

mvn clean package

Maven은 현재 폴더에 kafka-connect-protobuf-converter-<version>-jar-with-dependencies.jar라는 파일을 빌드해요. 이것이 변환기 JAR 파일이에요.

  • 컴파일된 kafka-connect-protobuf-converter-<version>-jar-with-dependencies.jar 파일을 Kafka 패키지 버전용 디렉토리에 복사하세요:
    • Confluent: <confluent_dir>/share/java/kafka-serde-tools
    • Apache Kafka: <apache_kafka_dir>/libs

2단계: .proto 파일 컴파일

메시지를 정의하는 protobuf .proto 파일을 java 파일로 컴파일하세요.

예를 들어 메시지가 sensor.proto라는 파일에 정의되어 있다고 가정하세요. 터미널 창에서 프로토콜 버퍼 파일을 컴파일하는 다음 명령을 실행하세요. 애플리케이션 소스 코드의 소스 디렉토리, 대상 디렉토리(.java 파일용), .proto 파일 경로를 지정하세요:

protoc -I=$SRC_DIR --java_out=$DST_DIR $SRC_DIR/sensor.proto

샘플 .proto 파일은 여기에서 사용할 수 있어요: https://github.com/snowflakedb/snowflake-kafka-connector/blob/master/test/test_data/sensor.proto

이 명령은 지정된 대상 디렉토리에 SensorReadingImpl.java라는 파일을 생성해요. 자세한 내용은 Google 개발자 문서를 참고하세요.

3단계: SensorReadingImpl.java 파일 컴파일

2단계: .proto 파일 컴파일에서 생성된 SensorReadingImpl.java 파일을 protobuf 프로젝트 구조의 Project Object Model과 함께 컴파일하세요.

  • 텍스트 편집기에서 .pom 파일을 여세요.
  • 그 외에는 비어 있는 디렉토리를 다음과 같은 구조로 만드세요:
protobuf_folder
├── pom.xml
└── src
 └── main
     └── java
         └── com
             └── ..

여기서 src/main/java 아래의 디렉토리 구조는 .proto 파일(3번째 줄)의 패키지 이름을 반영해요.

  • 2단계에서 생성된 SensorReadingImpl.java 파일을 디렉토리 구조의 맨 아래 폴더로 복사하세요.
  • protobuf_folder 디렉토리 루트에 pom.xml이라는 파일을 만드세요.
  • 빈 pom.xml 파일을 텍스트 편집기에서 열고 다음 예제 프로젝트 모델을 파일에 복사해 수정하세요:
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
     xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
     xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>

<groupId><group_id></groupId>
<artifactId><artifact_id></artifactId>
<version><version></version>

<properties>
    <java.version><java_version></java.version>
    <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
</properties>

<dependencies>
    <dependency>
        <groupId>com.google.protobuf</groupId>
        <artifactId>protobuf-java</artifactId>
        <version>3.11.1</version>
    </dependency>
</dependencies>

<build>
    <plugins>
        <plugin>
            <groupId>org.apache.maven.plugins</groupId>
            <artifactId>maven-compiler-plugin</artifactId>
            <version>3.3</version>
            <configuration>
                <source>${java.version}</source>
                <target>${java.version}</target>
            </configuration>
        </plugin>
        <plugin>
            <artifactId>maven-assembly-plugin</artifactId>
            <version>3.1.0</version>
            <configuration>
                <descriptorRefs>
                    <descriptorRef>jar-with-dependencies</descriptorRef>
                </descriptorRefs>
            </configuration>
            <executions>
                <execution>
                    <id>make-assembly</id>
                    <phase>package</phase>
                    <goals>
                        <goal>single</goal>
                    </goals>
                </execution>
            </executions>
        </plugin>
    </plugins>
</build>
</project>

여기서:

  • <group_id>: .proto 파일에 지정된 패키지 이름의 Group ID 부분. 예를 들어 패키지 이름이 com.foo.bar.buz이면 group ID는 com.foo입니다.
  • <artifact_id>: 선택한 패키지의 Artifact ID. artifact ID는 임의로 선택할 수 있어요.
  • <version>: 선택한 패키지의 버전. 버전은 임의로 선택할 수 있어요.
  • <java_version>: 로컬 머신에 설치된 JRE(Java Runtime Environment) 버전.

예를 들어:

<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
     xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
     xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>

<groupId>com.snowflake</groupId>
<artifactId>kafka-test-protobuf</artifactId>
<version>1.0.0</version>

<properties>
    <java.version>1.8</java.version>
    <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
</properties>

<dependencies>
    <dependency>
        <groupId>com.google.protobuf</groupId>
        <artifactId>protobuf-java</artifactId>
        <version>3.11.1</version>
    </dependency>
</dependencies>

<build>
    <plugins>
        <plugin>
            <groupId>org.apache.maven.plugins</groupId>
            <artifactId>maven-compiler-plugin</artifactId>
            <version>3.3</version>
            <configuration>
                <source>${java.version}</source>
                <target>${java.version}</target>
            </configuration>
        </plugin>
        <plugin>
            <artifactId>maven-assembly-plugin</artifactId>
            <version>3.1.0</version>
            <configuration>
                <descriptorRefs>
                    <descriptorRef>jar-with-dependencies</descriptorRef>
                </descriptorRefs>
            </configuration>
            <executions>
                <execution>
                    <id>make-assembly</id>
                    <phase>package</phase>
                    <goals>
                        <goal>single</goal>
                    </goals>
                </execution>
            </executions>
        </plugin>
    </plugins>
</build>
</project>
  • 터미널 창에서 protobuf_folder 디렉토리 루트로 변경하고 다음 명령을 실행해 디렉토리의 파일에서 protobuf 데이터 JAR 파일을 컴파일하세요:
mvn clean package

Maven은 protobuf_folder/target 폴더에 <artifact_id>-<version>-jar-with-dependencies.jar라는 파일(예: kafka-test-protobuf-1.0.0-jar-with-dependencies.jar)을 생성해요.

  • 컴파일된 kafka-test-protobuf-1.0.0-jar-with-dependencies.jar 파일을 Kafka 패키지 버전용 디렉토리에 복사하세요:
    • Confluent: <confluent_dir>/share/java/kafka-serde-tools
    • Apache Kafka: 파일을 $CLASSPATH 환경 변수 디렉토리에 복사하세요.

4단계: Kafka 커넥터 구성

  • 텍스트 편집기에서 Kafka 구성 파일(예: <kafka_dir>/config/connect-distributed.properties)을 여세요.
  • 파일에 value.converter.protoClassName 속성을 추가하세요. 이 속성은 메시지를 역직렬화하는 데 사용할 프로토콜 버퍼 클래스를 지정해요(예: com.google.protobuf.Int32Value). (참고: 중첩 클래스는 $ 표기법으로 지정해야 해요(예: com.blueapron.connect.protobuf.NestedTestProtoOuterClass$NestedTestProto).)
{
 "name":"XYZCompanySensorData",
"config":{
  ..
  "value.converter.protoClassName":"com.snowflake.kafka.test.protobuf.SensorReadingImpl$SensorReading"
}
 }

일반적인 Kafka 커넥터 속성에 대한 내용은 Kafka 구성 속성을 참고하세요. 프로토콜 버퍼 클래스에 대한 자세한 내용은 이 문서 앞부분에서 참조한 Google 개발자 문서를 참고하세요.

  • 파일을 저장하세요.

더 알아보기 (Learn more)