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
相关产品推荐
相关产品推荐

