Kafka コネクタは、Protocol Buffers (Protobuf) 形式のデータの読み取りをサポートしています。
Protocol Buffers
Protocol Buffers (Protobuf) は、Google によって開発された、効率的で言語に中立な構造化データシリアル化形式です。JSON や XML と比較して、次のような大きな利点があります:
-
コンパクトなサイズ:シリアル化されたデータはよりコンパクトになり、ストレージ容量とネットワーク帯域幅を節約します。
-
高速性:シリアル化と逆シリアル化が高速であるため、高性能なアプリケーションに最適です。
-
構造化された定義:データ構造を
.protoファイルで定義するため、明確で保守しやすいインターフェイスを提供します。 -
クロス言語サポート:主要なプログラミング言語をサポートしており、異なるシステム間でのデータ交換を容易にします。
これらの利点により、Protobuf は高頻度の通信、マイクロサービス、リアルタイムコンピューティングなどのシナリオで広く使用されています。Kafka で推奨される効率的なデータ形式の 1 つです。
制限事項
Kafka コネクタは、Protocol Buffers バージョン 21.7 以前をサポートしています。
ステップ 1:Protobuf ファイルのコンパイル
-
order.proto という名前の Protobuf ファイルを作成します。
proto3
syntax = "proto3"; // 他の .proto ファイルでこのファイルをインポートするための論理パッケージ名。 package com.aliyun; // Java パッケージ名。指定しない場合、デフォルトで proto パッケージが使用されます。 option java_package = "com.aliyun"; // 複数のファイルにコンパイルするかどうかを指定します。true に設定することを推奨します // これにより、各メッセージが内部クラスではなく、個別の .java ファイルを生成するようになります。 option java_multiple_files = true; // Java の外部クラス名。java_multiple_files が true の場合、このクラスはファイル名などのメタデータを含み、 // メッセージクラスをラップしません。 option java_outer_classname = "OrderProtoBuf"; message Order { // Proto3 では optional/required キーワードが削除されています。 // 注意:プリミティブデータ型 (int, long, double) のデフォルトは 0、string のデフォルトは空文字列です。 // Java では hasOrderId() メソッドが生成されなくなるため、 // 「設定されていない」と「0 に設定されている」を区別できなくなります。 int32 orderId = 1; string orderName = 2; double orderPrice = 3; int64 orderDate = 4; }proto2
syntax = "proto2"; // proto のパッケージ名。 package com.aliyun; // Java パッケージ名。指定しない場合、デフォルトで proto パッケージが使用されます。 option java_package = "com.aliyun"; // 複数のファイルにコンパイルするかどうかを指定します。 option java_multiple_files = true; // Java のラッパークラス名。 option java_outer_classname = "OrderProtoBuf"; message Order { optional int32 orderId = 1; optional string orderName= 2; optional double orderPrice = 3; optional int64 orderDate = 4; } -
Protocol Buffers ツールを使用してソースコードを生成します。
空の Maven プロジェクトを作成し、Protobuf ファイルを src/main/proto ディレクトリに配置します。
ディレクトリの例
KafkaProtobuf ‒ src -main -java -proto -order.proto ‒ pom.xmlpom.xml
説明生成されたクラスが protobuf-java:3.21.7 と競合するのを防ぐため、バージョンを Flink の依存関係 (例:3.21.7) と合わせてください。
<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/maven-v4_0_0.xsd"> <modelVersion>4.0.0</modelVersion> <groupId>com.aliyun</groupId> <artifactId>KafkaProtobuf</artifactId> <version>1.0-SNAPSHOT</version> <dependencies> <dependency> <groupId>com.google.protobuf</groupId> <artifactId>protobuf-java</artifactId> <version>3.21.7</version> <!-- このバージョンは、コード生成に使用される Protobuf のバージョンと一致させる必要があります。 --> </dependency> <dependency> <groupId>org.apache.kafka</groupId> <artifactId>kafka-clients</artifactId> <version>3.3.1</version> <!-- ご利用の Kafka のバージョンに合わせて調整してください。 --> </dependency> <dependency> <groupId>com.github.javafaker</groupId> <artifactId>javafaker</artifactId> <version>1.0.2</version> </dependency> </dependencies> <build> <finalName>KafkaProtobuf</finalName> <plugins> <!-- Java コンパイラ --> <plugin> <groupId>org.apache.maven.plugins</groupId> <artifactId>maven-compiler-plugin</artifactId> <version>3.13.0</version> <configuration> <source>1.8</source> <target>1.8</target> </configuration> </plugin> <!-- maven-shade-plugin を使用して、必要なすべての依存関係を含む fat JAR を作成します。 --> <plugin> <groupId>org.apache.maven.plugins</groupId> <artifactId>maven-shade-plugin</artifactId> <version>3.5.3</version> </plugin> </plugins> </build> </project>java ディレクトリに、Order クラス、OrderOrBuilder インターフェイス、および OrderProtobuf 外部ラッパークラスの 3 つのクラスが生成されます。
# ターミナルで、プロジェクトのルートディレクトリ (pom.xml がある場所) に移動し、 # 次のコマンドを実行してソースコードを生成します。 protoc --java_out=src/main/java --proto_path=src/main/proto src/main/proto/order.proto -
シリアル化と逆シリアル化をテストします。
package com.aliyun; public class OrderTest { public static void main(String[] args) { // Order オブジェクトを作成し、そのフィールド値を設定します。 Order order = Order.newBuilder() .setOrderId(8513) .setOrderName("flink") .setOrderPrice(99.99) .setOrderDate(System.currentTimeMillis()) .build(); // バイト配列にシリアル化します。 byte[] serializedBytes = order.toByteArray(); System.out.println("Byte length after serialization: " + serializedBytes.length); // バイト配列を新しい Order オブジェクトに逆シリアル化します。 Order deserializedOrder; try { deserializedOrder = Order.parseFrom(serializedBytes); } catch (Exception e) { System.err.println("Deserialization failed: " + e.getMessage()); return; } System.out.println("Original object: \n" + order); // 逆シリアル化されたオブジェクトのフィールドが元のオブジェクトと一致することを確認します。 if (order.getOrderId() == deserializedOrder.getOrderId() && order.getOrderName().equals(deserializedOrder.getOrderName()) && order.getOrderPrice() == deserializedOrder.getOrderPrice() && order.getOrderDate() == deserializedOrder.getOrderDate()) { System.out.println("Serialization and deserialization test passed!"); } else { System.out.println("Serialization and deserialization test failed!"); } } }
ステップ 2:Kafka へのテストデータの書き込み
この例では、ApsaraMQ for Kafka を動作環境として使用します。
-
SSL ルート証明書をダウンロードします。この証明書は、SSL エンドポイントを使用して接続する場合に必要です。
-
ご利用のインスタンスのユーザー名とパスワードを使用します。
-
インスタンスでアクセス制御リスト (ACL) が有効になっていない場合、ApsaraMQ for Kafka コンソールの [インスタンス詳細] ページの [設定情報] セクションからデフォルトのユーザー名とパスワードを取得できます。
-
インスタンスでアクセス制御リスト (ACL) が有効になっている場合、SASL ユーザーが PLAIN タイプを使用し、メッセージを送受信する権限を持っていることを確認してください。詳細については、「ACL を使用したアクセス制御」をご参照ください。
-
package com.aliyun;
import org.apache.kafka.clients.CommonClientConfigs;
import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.clients.producer.RecordMetadata;
import org.apache.kafka.common.config.SaslConfigs;
import org.apache.kafka.common.config.SslConfigs;
import java.util.ArrayList;
import java.util.List;
import java.util.Properties;
import java.util.concurrent.Future;
import com.github.javafaker.Faker; // Faker ライブラリをインポートして、ランダムなテストデータセットを生成します。
public class ProtoBufToKafkaTest {
public static void main(String[] args) {
Properties props = new Properties();
// エンドポイントを設定します。トピックのエンドポイントは Kafka コンソールから取得します。
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "<bootstrap_servers>");
// アクセスプロトコルを SASL_SSL に設定します。
props.put(CommonClientConfigs.SECURITY_PROTOCOL_CONFIG, "SASL_SSL");
// SSL ルート証明書の絶対パスを設定します。このファイルは JAR にパッケージ化しないでください。
props.put(SslConfigs.SSL_TRUSTSTORE_LOCATION_CONFIG, "../only.4096.client.truststore.jks");
// ルート証明書トラストストアのパスワード。デフォルト値を使用します。
props.put(SslConfigs.SSL_TRUSTSTORE_PASSWORD_CONFIG, "KafkaOnsClient");
// SASL 認証方式。デフォルト値を使用します。
props.put(SaslConfigs.SASL_MECHANISM, "PLAIN");
// ApsaraMQ for Kafka のメッセージキーと値のシリアル化メソッド。
props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer");
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.ByteArraySerializer");
// リクエストの最大待機時間 (ミリ秒)。
props.put(ProducerConfig.MAX_BLOCK_MS_CONFIG, 30 * 1000);
// クライアントの再試行回数。
props.put(ProducerConfig.RETRIES_CONFIG, 5);
// クライアントの再試行間隔 (ミリ秒)。
props.put(ProducerConfig.RETRY_BACKOFF_MS_CONFIG, 3000);
// 値を空文字列に設定して、ホスト名の検証を無効にします。
props.put(SslConfigs.SSL_ENDPOINT_IDENTIFICATION_ALGORITHM_CONFIG, "");
props.put("sasl.jaas.config", "org.apache.kafka.common.security.plain.PlainLoginModule required username=\"aliyun_flink\" password=\"123456\";");
// Producer オブジェクトを構築します。このオブジェクトはスレッドセーフです。
// 通常、プロセスごとに 1 つの Producer オブジェクトで十分です。
// パフォーマンスを向上させるために、より多くのオブジェクトを作成できますが、5 つを超えないようにしてください。
KafkaProducer<String, byte[]> producer = new KafkaProducer<>(props);
String topic = "test";
// 3 つのメッセージを格納するリストを作成します。
List<ProducerRecord<String, byte[]>> messages = new ArrayList<>();
for (int i = 0; i < 3; i++) {
byte[] value = getProtoTestData();
ProducerRecord<String, byte[]> kafkaMessage = new ProducerRecord<>(topic, value);
messages.add(kafkaMessage);
}
try {
// メッセージをバッチで送信します。
List<Future<RecordMetadata>> futures = new ArrayList<>();
for (ProducerRecord<String, byte[]> message : messages) {
Future<RecordMetadata> metadataFuture = producer.send(message);
futures.add(metadataFuture);
}
producer.flush();
// Future オブジェクトの結果を同期的に取得します。
for (Future<RecordMetadata> future : futures) {
try {
RecordMetadata recordMetadata = future.get();
System.out.println("Produce ok:" + recordMetadata.toString());
} catch (Throwable t) {
t.printStackTrace();
}
}
} catch (Exception e) {
// 再試行してもメッセージの送信に失敗する場合、ビジネスロジックでこのエラーを処理する必要があります。
System.out.println("error occurred");
e.printStackTrace();
}
}
private static byte[] getProtoTestData() {
// Faker を使用してランダムデータを生成します。
Faker faker = new Faker();
int orderId = faker.number().numberBetween(1000, 9999); // ランダムな注文 ID を生成します。
String orderName = faker.commerce().productName(); // ランダムな注文名を生成します。
double orderPrice = faker.number().randomDouble(2, 10, 1000); // ランダムな注文価格を生成します。
long orderDate = System.currentTimeMillis(); // 現在の時刻を注文日として使用します。
// 定義されたデータ構造に基づいてオブジェクトを作成します。
Order order = Order.newBuilder()
.setOrderId(orderId)
.setOrderName(orderName)
.setOrderPrice(orderPrice)
.setOrderDate(orderDate)
.build();
// データをシリアル化:オブジェクトデータをバイト配列に変換します。
return order.toByteArray();
}
}
-
テストコードを実行して、3 つの Protobuf 形式のメッセージを Kafka の
testトピックに書き込みます。Produce ok:test-1@3 Produce ok:test-1@4 Produce ok:test-1@5
ステップ 3:アーティファクトのビルドとアップロード
コンパイルおよびパッケージ化された KafkaProtobuf.jar ファイルをアップロードします。
左側のナビゲーションウィンドウで、[アーティファクト管理] をクリックします。[アーティファクト] タブで、[アーティファクトのアップロード] をクリックします。
組み込みの Protobuf データ形式は、Ververica Runtime (VVR) 8.0.9 以降でのみ利用可能です。それ以前のバージョンを使用する場合は、flink-protobuf-1.17.2.jar の依存関係を追加する必要があります。
ステップ 4:Flink SQL によるデータの読み取り
-
参考として、次の SQL の例をご参照ください。
protobuf.message-class-nameパラメーターを、Protobuf メッセージの完全修飾クラス名に設定します。protobufパラメーターの詳細については、「Flink-Protobuf」をご参照ください。CREATE TEMPORARY TABLE KafkaSource ( orderId INT, orderName STRING, orderPrice DOUBLE, orderDate BIGINT ) WITH ( 'connector' = 'kafka', 'topic' = 'test', 'properties.group.id' = 'my-group', -- コンシューマーグループの ID。 'properties.bootstrap.servers' = '<bootstrap_servers>', -- Kafka ブローカーのアドレスを入力します。 'format' = 'protobuf', -- value 部分のデータ形式。 'protobuf.message-class-name' = 'com.aliyun.Order', -- メッセージ本文のメッセージクラス。 'scan.startup.mode' = 'earliest-offset' -- Kafka パーティションの最小オフセットから読み取ります。 ); CREATE TEMPORARY TABLE KafkaSink ( orderId INT, orderName STRING, orderPrice DOUBLE, orderDate BIGINT ) WITH ( 'connector' = 'print' ); INSERT INTO KafkaSink SELECT * FROM KafkaSource ; -
追加の依存関係を参照します。
[詳細設定] パネルの [追加の依存関係] セクションで
KafkaProtobuf.jarをアップロードします。CREATE TEMPORARY TABLE文を使用して、フィールドorderId INT、orderName STRING、orderPrice DOUBLE、およびorderDate BIGINTを持つKafkaSourceテーブルを定義します。WITH 句で、connectorをkafkaに、formatをprotobufに、protobuf.message-class-nameをcom.aliyun.Orderに設定します。ジョブが正常に実行されると、消費されたデータが KafkaSink の結果タブに表示されます。 -
SQL コードをデバッグします。
[デバッグ] をクリックし、次の SQL 文を実行して Kafka Protobuf ソーステーブルからデータを読み取ります。
CREATE TEMPORARY TABLE KafkaSource ( orderId INT, orderName STRING, orderPrice DOUBLE, orderDate BIGINT ) WITH ( 'connector' = 'kafka', 'topic' = 'test', 'properties.group.id' = 'my-group', 'properties.bootstrap.servers' = 'alikafka-serverless-cn-xxx', 'format' = 'protobuf', 'protobuf.message-class-name' = 'com.aliyun.Order', 'scan.startup.mode' = 'earliest-offset' ); -
ジョブの実行開始後に出力を表示します。
ジョブ詳細ページで、[ジョブログ] タブに移動し、次に [実行中のタスクマネージャー] サブタブに移動します。ご利用のタスクマネージャー ID をクリックし、[Stdout] タブを選択します。Stdout の出力で、
+Iで始まるデータレコード (例:[6066, Small Marble Car, 150.83, 1745561043465]) を確認できます。これは、Flink SQL が Kafka から Protobuf 形式のデータを正常に読み取ったことを示します。
よくある質問
-
上流のソースから読み取り、下流の別の Kafka トピックに書き込んだ後、ログに多数の
CORRUPT_MESSAGE警告が表示されるのはなぜですか?原因:ApsaraMQ for Kafka では、Professional (High-Write) Edition 以外のインスタンスでローカルストレージを使用するトピックは、べき等な書き込みまたはトランザクション書き込みをサポートしていません。その結果、Kafka 結果テーブルが提供する exactly-once セマンティクスを使用できません。
解決策:結果テーブルに設定プロパティ
properties.enable.idempotence=falseを追加して、べき等な書き込み機能を無効にします。 -
実行時にジョブログで
NoClassDefFoundErrorが報告されるのはなぜですか?原因:アップロードされた protobuf-java JAR ファイルのバージョンが、Protocol Buffers コンパイラで使用されているバージョンと一致していません。
解決策:追加の依存関係のバージョンが一貫していること、欠落しているファイルがないこと、およびプロジェクトが正しくコンパイルおよびパッケージ化されていることを確認してください。
-
ジョブの検証が次のエラーで失敗するのはなぜですか:
Could not find any factory for identifier 'protobuf' that implements one of 'org.apache.flink.table.factories.EncodingFormatFactory'?原因:組み込みの Protobuf データ形式は、Ververica Runtime (VVR) 8.0.9 以降でのみサポートされています。
解決策:
flink-protobufの依存関係が追加されているかどうかを確認してください。