如何还原被UTF-8编码损坏的Protobuf二进制字节数组?
- 生产者将Protobuf消息序列化为二进制字节数组发送
- 配置错误的Kafka集群将该二进制数据反序列化为字符串
- 集群再将字符串序列化后发送给消费者
- 消费者期望接收二进制字节数组,但收到的是被UTF-8编码损坏的数据
复现用Proto文件
syntax = "proto3"; import "google/protobuf/wrappers.proto"; import "google/protobuf/timestamp.proto"; option java_package = "com.mycompany.proto"; option java_multiple_files = true; package com.mycompany; enum MessageType { NOT_SET = 0; TYPE_A = 1; TYPE_B = 2; } message MyMessagePart { string someValue = 1; } message MyMessage { // Numeric (integer) variable int32 myNumber = 1; // Text value string myText = 2; // Enum value MessageType mType = 3; // Message parts repeated MyMessagePart messagePart = 4; // Uint32 value google.protobuf.UInt32Value uint32Value = 5; // Timestamp google.protobuf.Timestamp timestamp = 6; }
JUnit复现测试代码
public class EncodingTest { @Test public void dealWithCorruptedBinaryData() throws InvalidProtocolBufferException { // 1. 创建Protobuf消息 final MyMessage msg = MyMessage.newBuilder() .setMyNumber(42) .setMyText("Hello") .setMType(MessageType.TYPE_A) .setUint32Value(UInt32Value.newBuilder() .setValue(2067) .build()) .addMessagePart(MyMessagePart.newBuilder() .setSomeValue("message part value") .build()) .build(); // 2. 转换为字节数组 final byte[] bytesSentByProducer = msg.toByteArray(); // 3. 模拟配置错误的Kafka反序列化二进制为字符串 final StringDeserializer deserializer = new StringDeserializer(); final String dataReceivedInsideMisconfiguredKafka = deserializer.deserialize("inputTopic", bytesSentByProducer); // 4. 模拟Kafka将字符串重新序列化发送给消费者 final StringSerializer serializer = new StringSerializer(); final byte[] dataSentToConsumer = serializer.serialize("outputTopic", dataReceivedInsideMisconfiguredKafka); // 损坏后的字节数组无法解析为Protobuf消息 final MyMessage receivedMessage = MyMessage.parseFrom(dataSentToConsumer); } }
解析异常信息
尝试解析损坏后的字节数组时,抛出以下异常:
com.google.protobuf.InvalidProtocolBufferException: While parsing a protocol message, the input ended unexpectedly in the middle of a field. This could mean either that the input has been truncated or that an embedded message misreported its own length. at com.google.protobuf.InvalidProtocolBufferException.truncatedMessage(InvalidProtocolBufferException.java:107) at com.google.protobuf.CodedInputStream$ArrayDecoder.readRawByte(CodedInputStream.java:1245) at com.google.protobuf.CodedInputStream$ArrayDecoder.readRawVarint64SlowPath(CodedInputStream.java:1130) at com.google.protobuf.CodedInputStream$ArrayDecoder.readRawVarint32(CodedInputStream.java:1024) at com.google.protobuf.CodedInputStream$ArrayDecoder.readUInt32(CodedInputStream.java:954) at com.google.protobuf.UInt32Value.<init>(UInt32Value.java:58) at com.google.protobuf.UInt32Value.<init>(UInt32Value.java:14)
注:使用未损坏的
bytesSentByProducer可以正常解析为Protobuf消息。
核心问题
- 是否可以将
dataSentToConsumer转换回原始的bytesSentByProducer? - 若可行,在仅能控制消费者的情况下,如何修复该问题?如何撤销Kafka集群中发生的UTF-8编码操作?
注:正确解决方案是修正Kafka集群配置,但因流程限制无法实施。
已尝试的修复方法
方法1:使用Charset编解码逆操作
private byte[] convertToOriginalBytes(final byte[] bytesAfter) throws CharacterCodingException { final Charset charset = StandardCharsets.UTF_8; final CharsetDecoder decoder = charset.newDecoder(); final CharsetEncoder encoder = charset.newEncoder(); final ByteBuffer byteBuffer = ByteBuffer.wrap(bytesAfter); final CharBuffer charBuffer = CharBuffer.allocate(bytesAfter.length); final CoderResult result = decoder.decode(byteBuffer, charBuffer, true); result.throwException(); final ByteBuffer reversedByteBuffer = encoder.encode(charBuffer); final byte[] reversedBytes = new byte[reversedByteBuffer.remaining()]; reversedByteBuffer.get(reversedBytes); return reversedBytes; }
执行后抛出异常:
java.nio.BufferUnderflowException at java.base/java.nio.charset.CoderResult.throwException(CoderResult.java:272) at com.mycompany.EncodingTest.convertToOriginalBytes(EncodingTest.java:67) at com.mycompany.EncodingTest.dealWithCorruptedBinaryData(EncodingTest.java:54)
方法2:基于UTF-8字节模式的位操作推测
已知UTF-8字节模式:
0xxxxxxx:单字节字符110xxxxx 10xxxxxx:双字节字符等
推测StringDeserializer/StringSerializer会修改二进制数据以符合UTF-8规则,若转换可逆,可通过位操作还原原始消息,但尚未找到可行实现。
解决方案
1. 数据恢复可能性判断
Kafka默认配置的StringDeserializer会将无效UTF-8字节替换为U+FFFD(替换字符),这一步是不可逆的——多个不同的无效字节序列都会被替换为同一个U+FFFD,无法区分原始字节值。因此,对于包含无效UTF-8字节的Protobuf二进制数据(绝大多数场景),无法完全恢复原始的bytesSentByProducer。
仅当原始Protobuf字节恰好全部是合法UTF-8序列(极罕见)时,才能通过UTF-8解码再编码的方式无损恢复,但这种情况几乎不会出现(Protobuf的Varint编码会生成0x80及以上的单字节,属于无效UTF-8序列)。
2. 消费者端可行修复方案
方案A:部分恢复合法UTF-8字节段
如果业务可以接受部分数据丢失,或能确定不受损坏影响的字段,可以尝试以下代码恢复原本合法的UTF-8字节:
// 将损坏的字节数组按UTF-8解码为字符串,再用ISO-8859-1编码回字节数组 // 注意:无效字节会被替换为0xEF 0xBF 0xBD,无法还原原始值 byte[] recoveredBytes = new String(dataSentToConsumer, StandardCharsets.UTF_8).getBytes(StandardCharsets.ISO_8859_1);
但这种方式无法修复Protobuf解析异常,仅适用于特定场景。
方案B:规避损坏流程的替代方案
若无法修改Kafka配置,可尝试协调生产者端临时调整发送方式:将二进制字节数组用Base64编码为字符串后发送,这样Kafka的StringDeserializer/Serializer不会损坏数据,消费者端再解码Base64得到原始字节数组。但此方案需要生产者配合修改。
3. 结论
对于绝大多数Protobuf场景,经过Kafka的UTF-8编解码损坏后,无法完全恢复原始数据。消费者端只能做有限的部分恢复,或等待Kafka集群配置修正。
内容的提问来源于stack exchange,提问作者PravlesRedneckoff

