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

Kafka JSON消息中Base64编码的Protobuf payload如何正确反序列化

问题原因与解决方案

错误根因

  • 序列化协议不匹配:ObjectInputStream是Java原生序列化体系的工具,仅能解析ObjectOutputStream输出的、符合Java序列化协议的字节流。而你拿到的payload字段Base64解码后的字节是Protobuf自定义的序列化格式,和Java序列化协议完全不兼容,缺少Java序列化要求的固定流头标识,因此抛出StreamCorruptedException异常。
  • 代码语法错误:你代码中byte bytes = Base64.getDecoder().decode(json.getString("payload"));存在两处错误:一是Base64.Decoder.decode返回值为byte[]类型,不是单个byte;二是你定义的JSON对象变量名为jsonObj,调用时误写为json。

正确实现方案

Protobuf编译生成的消息类自带反序列化方法,直接调用静态方法parseFrom传入Base64解码后的字节数组即可完成反序列化,无需使用Java原生序列化的流工具。
修改后的消费端核心代码如下:

// 需额外导入Protobuf反序列化异常类
import com.google.protobuf.InvalidProtocolBufferException;

// 循环拉取消息的逻辑替换为以下实现
while(true){  
    ConsumerRecords<String,String> records=consumer.poll(Duration.ofMillis(100));  
    for(ConsumerRecord<String,String> record: records){
        try {
            JSONObject jsonObj = new JSONObject(record.value());
            // 1. 修正变量类型,取出payload并做Base64解码
            byte[] payloadBytes = Base64.getDecoder().decode(jsonObj.getString("payload"));
            // 2. 直接调用Protobuf生成类的parseFrom方法反序列化
            ExampleProtocolBuffer.messageData message = ExampleProtocolBuffer.messageData.parseFrom(payloadBytes);
            logger.info("Decoded FirstName: " + message.getName().getFirstName());
        } catch (JSONException | IllegalArgumentException | InvalidProtocolBufferException e) {
            // 分别处理JSON解析异常、Base64解码异常、Protobuf反序列化异常
            e.printStackTrace();
        }
    }      
} 

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.05 06:12:03