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

如何使用ZIO Kafka读取Kafka Topic并反序列化为Java Protobuf类

解决ZIO Kafka读取Protobuf数据转自定义Java类的ClassCastException问题

核心原因

ZIO Kafka默认的Protobuf反序列化器会把消息解析为DynamicMessage,而非你通过proto文件生成的proto.Data静态类,直接调用getId()或强转自然会抛出类型转换异常。

解决方案步骤

1. 实现自定义Protobuf反序列化器

基于生成的proto.Data类,编写ZIO Kafka兼容的Deserializer,明确指定用静态类解析消息:

import com.google.protobuf.InvalidProtocolBufferException;
import zio.kafka.serde.Deserializer;
import zio.Chunk;

public class DataDeserializer implements Deserializer<proto.Data> {
    @Override
    public proto.Data deserialize(String topic, Chunk<Byte> bytes) {
        try {
            return proto.Data.parseFrom(bytes.toArray());
        } catch (InvalidProtocolBufferException e) {
            throw new RuntimeException("Failed to deserialize proto.Data", e);
        }
    }

    @Override
    public Deserializer<proto.Data> withSchemaRegistry(String url) {
        return this; // 无需Schema Registry时直接返回自身,需要则扩展逻辑
    }
}

2. 在消费者配置中指定自定义反序列化器

创建ZIO Kafka消费者时,将值的反序列化器替换为你自定义的实现,替代默认的Protobuf反序列化器:

import zio.kafka.consumer.Consumer;
import zio.kafka.consumer.ConsumerSettings;
import zio.kafka.serde.Serde;
import zio.ZIO;

// 配置消费者参数
ConsumerSettings settings = ConsumerSettings.builder()
    .bootstrapServers("localhost:9092")
    .groupId("data-consumer-group")
    .valueSerde(Serde.of(new DataDeserializer())) // 指定自定义反序列化器
    .build();

// 消费并过滤符合条件的记录
Consumer.consumeWith(settings, Subscription.topics("your-target-topic"), (key, value) -> {
    if (value.getId().equals("1")) { // 此时value为proto.Data类型,可直接调用方法
        System.out.println("Filtered valid record: " + value);
    }
    return ZIO.unit;
}).runDrain();

3. (可选)集成Schema Registry做Schema校验

如果需要校验消息Schema与本地类的一致性,可扩展反序列化器集成Schema Registry逻辑:

import io.confluent.kafka.schemaregistry.client.SchemaRegistryClient;
import io.confluent.kafka.schemaregistry.client.CachedSchemaRegistryClient;

public class SchemaValidatingDataDeserializer implements Deserializer<proto.Data> {
    private final SchemaRegistryClient schemaClient;

    public SchemaValidatingDataDeserializer(String registryUrl) {
        this.schemaClient = new CachedSchemaRegistryClient(registryUrl, 100);
    }

    @Override
    public proto.Data deserialize(String topic, Chunk<Byte> bytes) {
        // 从消息中提取Schema ID,校验是否与本地proto.Data的Schema匹配
        // 具体提取逻辑可参考Confluent Schema Registry官方文档
        try {
            return proto.Data.parseFrom(bytes.toArray());
        } catch (InvalidProtocolBufferException e) {
            throw new RuntimeException("Schema mismatch or deserialization failed", e);
        }
    }

    @Override
    public Deserializer<proto.Data> withSchemaRegistry(String url) {
        return new SchemaValidatingDataDeserializer(url);
    }
}

关键注意事项

  • 确保本地生成的proto.Data类的Schema与Kafka消息携带的Schema完全一致,否则会出现反序列化失败或校验不通过的问题。
  • 避免使用ZIO Kafka默认的Protobuf反序列化器,它仅适用于动态解析Protobuf消息的场景,不兼容静态生成的Protobuf类。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.25 19:48:05