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

能否在Spring-Kafka中直接使用Protobuf?序列化问题咨询

Spring Boot + Kafka Protobuf 交互(Confluent序列化/反序列化问题解决)

一、必备依赖配置

Maven 依赖

确保引入Spring Kafka、Confluent Protobuf序列化器及Protobuf核心库:

<dependencies>
    <!-- Spring Kafka -->
    <dependency>
        <groupId>org.springframework.kafka</groupId>
        <artifactId>spring-kafka</artifactId>
    </dependency>
    <!-- Confluent Protobuf 序列化器 -->
    <dependency>
        <groupId>io.confluent</groupId>
        <artifactId>kafka-protobuf-serializer</artifactId>
        <version>7.4.0</version> <!-- 与Kafka版本匹配 -->
    </dependency>
    <!-- Protobuf 核心库 -->
    <dependency>
        <groupId>com.google.protobuf</groupId>
        <artifactId>protobuf-java</artifactId>
        <version>3.24.4</version>
    </dependency>
</dependencies>

Maven Protobuf 编译插件

用于自动编译.proto文件生成Java类:

<build>
    <extensions>
        <extension>
            <groupId>kr.motd.maven</groupId>
            <artifactId>os-maven-plugin</artifactId>
            <version>1.7.1</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.24.4:exe:${os.detected.classifier}</protocArtifact>
                <pluginId>java</pluginId>
            </configuration>
            <executions>
                <execution>
                    <goals>
                        <goal>compile</goal>
                    </goals>
                </execution>
            </executions>
        </plugin>
    </plugins>
</build>

二、Protobuf 消息定义

使用proto3语法定义消息,确保字段编号唯一且类型明确:

syntax = "proto3";

package com.example.kafka.protobuf;

message UserEvent {
  int64 id = 1;
  string username = 2;
  string email = 3;
  int64 timestamp = 4;
}

编译后生成的Java类会存放在target/generated-sources/protobuf/java目录,需将该目录标记为IDE的生成源码根目录。

三、生产者配置

YAML 配置

spring:
  kafka:
    bootstrap-servers: localhost:9092
    producer:
      key-serializer: org.apache.kafka.common.serialization.StringSerializer
      value-serializer: io.confluent.kafka.serializers.protobuf.KafkaProtobufSerializer
      properties:
        schema.registry.url: http://localhost:8081 # 必须配置Confluent Schema Registry地址
        specific.protobuf.value.type: com.example.kafka.protobuf.UserEvent # 指定具体消息类型

Java 代码配置(可选)

@Configuration
public class KafkaProducerConfig {

    @Bean
    public ProducerFactory<String, UserEvent> producerFactory() {
        Map<String, Object> configProps = new HashMap<>();
        configProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
        configProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
        configProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, KafkaProtobufSerializer.class);
        configProps.put(KafkaProtobufSerializerConfig.SCHEMA_REGISTRY_URL_CONFIG, "http://localhost:8081");
        configProps.put(KafkaProtobufSerializerConfig.SPECIFIC_PROTOBUF_VALUE_TYPE, UserEvent.class.getName());
        return new DefaultKafkaProducerFactory<>(configProps);
    }

    @Bean
    public KafkaTemplate<String, UserEvent> kafkaTemplate() {
        return new KafkaTemplate<>(producerFactory());
    }
}

四、消费者配置

YAML 配置

spring:
  kafka:
    consumer:
      group-id: protobuf-consumer-group
      key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
      value-deserializer: io.confluent.kafka.serializers.protobuf.KafkaProtobufDeserializer
      properties:
        schema.registry.url: http://localhost:8081
        specific.protobuf.value.type: com.example.kafka.protobuf.UserEvent # 必须与生产者指定类型一致
        auto.offset.reset: earliest

Java 代码配置(可选)

@Configuration
public class KafkaConsumerConfig {

    @Bean
    public ConsumerFactory<String, UserEvent> consumerFactory() {
        Map<String, Object> props = new HashMap<>();
        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);
        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, KafkaProtobufDeserializer.class);
        props.put(KafkaProtobufDeserializerConfig.SCHEMA_REGISTRY_URL_CONFIG, "http://localhost:8081");
        props.put(KafkaProtobufDeserializerConfig.SPECIFIC_PROTOBUF_VALUE_TYPE, UserEvent.class.getName());
        props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
        return new DefaultKafkaConsumerFactory<>(props);
    }

    @Bean
    public ConcurrentKafkaListenerContainerFactory<String, UserEvent> kafkaListenerContainerFactory() {
        ConcurrentKafkaListenerContainerFactory<String, UserEvent> factory = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(consumerFactory());
        return factory;
    }
}

消费者监听示例

@Component
public class UserEventConsumer {

    @KafkaListener(topics = "user-events", groupId = "protobuf-consumer-group")
    public void consume(UserEvent event) {
        System.out.println("Received event: ID=" + event.getId() + ", Username=" + event.getUsername());
    }
}

五、常见解析错误排查

  • Schema Registry 不可达:确保schema.registry.url配置正确,且Schema Registry服务正常运行,Confluent序列化器依赖该服务存储/获取Protobuf schema。
  • 未指定具体消息类型:必须配置specific.protobuf.value.type,否则反序列化器无法确定目标类型,导致解析失败。
  • Protobuf 版本不兼容:生产者、消费者、Schema Registry使用的Protobuf版本需统一(如均为3.x),避免语法或序列化格式差异。
  • 消息格式不匹配:生产者必须使用Confluent的KafkaProtobufSerializer发送消息,不能用原生Protobuf序列化后直接发送,否则消费者无法解析。
  • Schema 兼容性问题:更新消息schema时需保证向前兼容(如添加可选字段),否则旧消费者无法解析新消息,反之亦然。
  • 生成类未被IDE识别:手动将target/generated-sources/protobuf/java标记为生成源码根目录,避免编译错误。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 23:43:21