Confluent 7.2.2中Protobuf特定消息SerDes使用异常求助
基于Confluent Platform 7.2.2的Kafka Streams Protobuf特定类SerDes完整实现
核心问题解析
- 未调用
configure时,KafkaProtobufSerde默认会将消息反序列化为DynamicMessage,强转为自定义AssetKey/AssetConfig类必然抛出类型转换异常。 - 调用
configure时触发schema无效异常,大概率是参数名错误(混淆了key/value的特定类参数)、生成的Protobuf类与Schema Registry中存储的schema不匹配,或是依赖版本不一致导致的。
完整实现步骤
1. 依赖配置(Maven)
确保所有Confluent相关依赖版本统一为7.2.2,同时引入Protobuf Java依赖:
<dependencies> <dependency> <groupId>org.apache.kafka</groupId> <artifactId>kafka-streams</artifactId> <version>7.2.2-ce</version> </dependency> <dependency> <groupId>io.confluent</groupId> <artifactId>kafka-streams-protobuf-serde</artifactId> <version>7.2.2</version> </dependency> <dependency> <groupId>io.confluent</groupId> <artifactId>kafka-protobuf-serializer</artifactId> <version>7.2.2</version> </dependency> <dependency> <groupId>io.confluent</groupId> <artifactId>kafka-schema-registry-client</artifactId> <version>7.2.2</version> </dependency> <dependency> <groupId>com.google.protobuf</groupId> <artifactId>protobuf-java</artifactId> <version>3.21.12</version> </dependency> </dependencies> <!-- Protobuf代码生成插件 --> <build> <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>${project.build.directory}/generated-sources/protobuf/java</outputDirectory> </configuration> <executions> <execution> <goals> <goal>compile</goal> </goals> </execution> </executions> </plugin> </plugins> </build>
2. Protobuf定义文件
创建src/main/proto/asset.proto,定义AssetKey和AssetConfig消息:
syntax = "proto3"; package com.example.protobuf; option java_multiple_files = true; option java_package = "com.example.protobuf"; option java_outer_classname = "AssetProto"; message AssetKey { string asset_id = 1; string asset_type = 2; } message AssetConfig { string config_key = 1; string config_value = 2; int64 update_time = 3; }
执行mvn compile后,会自动生成对应的Java类到target/generated-sources/protobuf/java目录。
3. Kafka Streams应用代码
package com.example; import com.example.protobuf.AssetConfig; import com.example.protobuf.AssetKey; import io.confluent.kafka.schemaregistry.client.CachedSchemaRegistryClient; import io.confluent.kafka.schemaregistry.client.SchemaRegistryClient; import io.confluent.kafka.serializers.protobuf.KafkaProtobufSerde; import org.apache.kafka.common.serialization.Serde; import org.apache.kafka.streams.KafkaStreams; import org.apache.kafka.streams.StreamsBuilder; import org.apache.kafka.streams.StreamsConfig; import org.apache.kafka.streams.Topology; import org.apache.kafka.streams.kstream.Consumed; import org.apache.kafka.streams.kstream.KStream; import org.apache.kafka.streams.kstream.Produced; import java.util.HashMap; import java.util.Map; import java.util.Properties; public class AssetStreamsApp { public static void main(String[] args) { // 基础Streams配置 Properties streamsProps = new Properties(); streamsProps.put(StreamsConfig.APPLICATION_ID_CONFIG, "asset-config-streams"); streamsProps.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); String schemaRegistryUrl = "http://localhost:8081"; // 初始化Schema Registry客户端 SchemaRegistryClient schemaRegistryClient = new CachedSchemaRegistryClient( schemaRegistryUrl, 100 // 缓存schema的数量 ); // 配置Key的SerDe(AssetKey) Serde<AssetKey> assetKeySerde = new KafkaProtobufSerde<>(schemaRegistryClient); Map<String, String> keySerdeConfigs = new HashMap<>(); keySerdeConfigs.put("schema.registry.url", schemaRegistryUrl); keySerdeConfigs.put("specific.protobuf.key.type", AssetKey.class.getName()); assetKeySerde.configure(keySerdeConfigs, true); // 第二个参数为true标识这是Key的SerDe // 配置Value的SerDe(AssetConfig) Serde<AssetConfig> assetConfigSerde = new KafkaProtobufSerde<>(schemaRegistryClient); Map<String, String> valueSerdeConfigs = new HashMap<>(); valueSerdeConfigs.put("schema.registry.url", schemaRegistryUrl); valueSerdeConfigs.put("specific.protobuf.value.type", AssetConfig.class.getName()); assetConfigSerde.configure(valueSerdeConfigs, false); // 第二个参数为false标识这是Value的SerDe // 构建拓扑 StreamsBuilder builder = new StreamsBuilder(); KStream<AssetKey, AssetConfig> assetStream = builder.stream( "asset-config-topic", Consumed.with(assetKeySerde, assetConfigSerde) ); // 示例处理逻辑:打印接收到的消息 assetStream.foreach((key, config) -> { System.out.printf("Asset ID: %s, Type: %s | Config: %s=%s (Updated: %d)%n", key.getAssetId(), key.getAssetType(), config.getConfigKey(), config.getConfigValue(), config.getUpdateTime()); }); // 输出处理后的消息到新主题 assetStream.to("processed-asset-config-topic", Produced.with(assetKeySerde, assetConfigSerde)); // 启动Streams应用 Topology topology = builder.build(); KafkaStreams streams = new KafkaStreams(topology, streamsProps); // 注册关闭钩子,优雅停止应用 Runtime.getRuntime().addShutdownHook(new Thread(streams::close)); streams.start(); } }
关键注意事项
- 参数名区分:Key的特定类参数是
specific.protobuf.key.type,Value的是specific.protobuf.value.type,不可混用。 - Schema一致性:生产者发送消息时,必须使用与Streams应用相同的Protobuf类生成消息,确保Schema Registry中存储的schema与本地类的descriptor完全匹配。
- 版本统一:所有Confluent生态的依赖版本必须与Confluent Platform版本(7.2.2)保持一致,避免序列化/反序列化兼容性问题。
- configure参数:
configure方法的第二个参数isKey必须正确设置,Key设为true,Value设为false,否则会导致schema处理逻辑异常。
内容的提问来源于stack exchange,提问作者devin
相关产品推荐
相关产品推荐

