すべてのプロダクト
Search
ドキュメントセンター

Realtime Compute for Apache Flink:Kafka コネクタを使用した Protobuf データの処理

最終更新日:Jun 22, 2026

Kafka コネクタは、Protocol Buffers (Protobuf) 形式のデータの読み取りをサポートしています。

Protocol Buffers

Protocol Buffers (Protobuf) は、Google によって開発された、効率的で言語に中立な構造化データシリアル化形式です。JSON や XML と比較して、次のような大きな利点があります:

  • コンパクトなサイズ:シリアル化されたデータはよりコンパクトになり、ストレージ容量とネットワーク帯域幅を節約します。

  • 高速性:シリアル化と逆シリアル化が高速であるため、高性能なアプリケーションに最適です。

  • 構造化された定義:データ構造を .proto ファイルで定義するため、明確で保守しやすいインターフェイスを提供します。

  • クロス言語サポート:主要なプログラミング言語をサポートしており、異なるシステム間でのデータ交換を容易にします。

これらの利点により、Protobuf は高頻度の通信、マイクロサービス、リアルタイムコンピューティングなどのシナリオで広く使用されています。Kafka で推奨される効率的なデータ形式の 1 つです。

制限事項

Kafka コネクタは、Protocol Buffers バージョン 21.7 以前をサポートしています。

ステップ 1:Protobuf ファイルのコンパイル

  1. 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;
    }
  2. Protocol Buffers ツールを使用してソースコードを生成します。

    空の Maven プロジェクトを作成し、Protobuf ファイルを src/main/proto ディレクトリに配置します。

    ディレクトリの例

    KafkaProtobuf
    ‒ src
      -main
        -java
        -proto
          -order.proto
    ‒ pom.xml

    pom.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
  3. シリアル化と逆シリアル化をテストします。

    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 を動作環境として使用します。

  1. SSL ルート証明書をダウンロードします。この証明書は、SSL エンドポイントを使用して接続する場合に必要です。

  2. ご利用のインスタンスのユーザー名パスワードを使用します。

    • インスタンスでアクセス制御リスト (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();
    }
}
  1. テストコードを実行して、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 によるデータの読み取り

  1. 参考として、次の 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
    ;
  2. 追加の依存関係を参照します。

    [詳細設定] パネルの [追加の依存関係] セクションで KafkaProtobuf.jar をアップロードします。CREATE TEMPORARY TABLE 文を使用して、フィールド orderId INTorderName STRINGorderPrice DOUBLE、および orderDate BIGINT を持つ KafkaSource テーブルを定義します。WITH 句で、connectorkafka に、formatprotobuf に、protobuf.message-class-namecom.aliyun.Order に設定します。ジョブが正常に実行されると、消費されたデータが KafkaSink の結果タブに表示されます。

  3. 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'
    );

  4. ジョブの実行開始後に出力を表示します。

    ジョブ詳細ページで、[ジョブログ] タブに移動し、次に [実行中のタスクマネージャー] サブタブに移動します。ご利用のタスクマネージャー 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 の依存関係が追加されているかどうかを確認してください。