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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 12:30:55