You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

求助:如何基于Protobuf格式编写Kafka Consumer?

Building a Kafka Consumer for Protobuf Messages

Got it, let's walk through building a Kafka Consumer that handles Protobuf messages step by step—this should get you up and running smoothly, since I know that repo you referenced didn't give you the clarity you needed.

Step 1: Define Your Protobuf Schema

First, you need a .proto file that defines your message structure. Let's use a simple example:

syntax = "proto3";

package com.example.proto;

message User {
  int32 id = 1;
  string name = 2;
  string email = 3;
}

Save this as src/main/proto/User.proto.

Step 2: Generate Protobuf Java Classes

You'll need to generate Java classes from your proto file. Using Maven? Add the Protobuf plugin to your pom.xml:

<build>
  <extensions>
    <extension>
      <groupId>kr.motd.maven</groupId>
      <artifactId>os-maven-plugin</artifactId>
      <version>1.7.0</version>
    </extension>
  </extensions>
  <plugins>
    <plugin>
      <groupId>org.xolstice.maven.plugins</groupId>
      <artifactId>protobuf-maven-plugin</artifactId>
      <version>0.6.1</version>
      <configuration>
        <protocArtifact>com.google.protobuf:protoc:3.21.12:exe:${os.detected.classifier}</protocArtifact>
        <pluginId>java</pluginId>
        <outputDirectory>target/generated-sources/protobuf/java</outputDirectory>
      </configuration>
      <executions>
        <execution>
          <goals>
            <goal>compile</goal>
          </goals>
        </execution>
      </executions>
    </plugin>
  </plugins>
</build>

Run mvn clean compile to generate the User.java class in the target directory.

Step 3: Add Kafka & Protobuf Dependencies

Include these dependencies in your pom.xml to work with Kafka and Protobuf serialization:

<dependencies>
  <!-- Kafka Clients -->
  <dependency>
    <groupId>org.apache.kafka</groupId>
    <artifactId>kafka-clients</artifactId>
    <version>3.4.0</version>
  </dependency>
  <!-- Protobuf Java -->
  <dependency>
    <groupId>com.google.protobuf</groupId>
    <artifactId>protobuf-java</artifactId>
    <version>3.21.12</version>
  </dependency>
  <!-- Confluent Protobuf Serializer (for Schema Registry integration) -->
  <dependency>
    <groupId>io.confluent</groupId>
    <artifactId>kafka-protobuf-serializer</artifactId>
    <version>7.4.0</version>
    <exclusions>
      <exclusion>
        <groupId>org.apache.kafka</groupId>
        <artifactId>kafka-clients</artifactId>
      </exclusion>
    </exclusions>
  </dependency>
</dependencies>

Step 4: Write the Kafka Consumer Code

Now let's create the consumer. We'll use Confluent's KafkaProtobufDeserializer which works seamlessly with Schema Registry (highly recommended for managing Protobuf schema versions):

import com.example.proto.User;
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.common.serialization.StringDeserializer;
import io.confluent.kafka.serializers.protobuf.KafkaProtobufDeserializer;
import io.confluent.kafka.serializers.protobuf.KafkaProtobufDeserializerConfig;

import java.time.Duration;
import java.util.Collections;
import java.util.Properties;

public class ProtobufKafkaConsumer {
    public static void main(String[] args) {
        Properties props = new Properties();
        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
        props.put(ConsumerConfig.GROUP_ID_CONFIG, "protobuf-consumer-group");
        props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
        
        // Configure Protobuf deserializer
        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, KafkaProtobufDeserializer.class.getName());
        props.put(KafkaProtobufDeserializerConfig.SCHEMA_REGISTRY_URL_CONFIG, "http://localhost:8081");
        props.put(KafkaProtobufDeserializerConfig.SPECIFIC_PROTOBUF_VALUE_TYPE, User.class.getName());

        try (KafkaConsumer<String, User> consumer = new KafkaConsumer<>(props)) {
            consumer.subscribe(Collections.singletonList("user-topic"));

            while (true) {
                ConsumerRecords<String, User> records = consumer.poll(Duration.ofMillis(100));
                records.forEach(record -> {
                    User user = record.value();
                    System.out.printf("Received user: ID=%d, Name=%s, Email=%s%n",
                            user.getId(), user.getName(), user.getEmail());
                    // Add your business logic here
                });
                consumer.commitAsync();
            }
        }
    }
}

Step 5: If You Don't Want to Use Schema Registry

If you're not using Schema Registry (not recommended for production, but okay for testing), you can create a custom deserializer that directly parses the byte array into your Protobuf object:

import com.example.proto.User;
import org.apache.kafka.common.serialization.Deserializer;

import java.util.Map;

public class CustomProtobufDeserializer implements Deserializer<User> {
    @Override
    public void configure(Map<String, ?> configs, boolean isKey) {}

    @Override
    public User deserialize(String topic, byte[] data) {
        if (data == null) return null;
        try {
            return User.parseFrom(data);
        } catch (Exception e) {
            throw new RuntimeException("Failed to deserialize Protobuf message", e);
        }
    }

    @Override
    public void close() {}
}

Then update your consumer config to use this custom deserializer:

props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, CustomProtobufDeserializer.class.getName());

Key Notes to Remember

  • Schema Compatibility: If using Schema Registry, ensure your Protobuf schema changes are backward/forward compatible to avoid deserialization failures.
  • Version Matching: Make sure your Kafka clients, Protobuf library, and Confluent serializer versions are compatible (check Confluent's version matrices for alignment).
  • Error Handling: Add try-catch blocks around message processing to handle cases where deserialization fails (e.g., corrupted messages, schema mismatches).

内容的提问来源于stack exchange,提问作者Akash Jain

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.27 09:59:10