Go如何读取二进制文件内多条protobuf记录并解码为JSON输出
问题描述
现有Kafka流转储文件,内部存储了数量未知的二进制格式Protobuf记录,需要将记录逐一解码转换为JSON格式输出到控制台。当前已尝试的代码仅能读取到第一条记录,代码如下:
package main import ( "encoding/json" "fmt" "github.com/golang/protobuf/proto" "io/ioutil" "parseRawDHCP/pb" ) func main() { file, err := ioutil.ReadFile("file") if err != nil { fmt.Printf("unable to read file %v", err) } msg := pb.Msg{} buffer := proto.NewBuffer(file) for { err := buffer.DecodeMessage(&msg) if err != nil { panic("unable to decode message") } marshalledStruct, err := json.Marshal(msg) if err != nil { panic("can't marshalledStruct the data from message") } if err == nil { fmt.Printf("message is: %v", marshalledStruct) continue } } }
解决方案
只能读取单条记录的核心原因有两个:
- Protobuf本身没有自边界特性,多条连续存储的Protobuf记录默认采用「长度前缀帧」格式存储,即每条消息前会追加一个varint编码的长度值标记当前消息的字节长度,你直接调用
DecodeMessage不会自动处理长度前缀,会把后续的长度字段当成消息内容解析,导致仅第一条可读 - 你复用了同一个
msg结构体且没有每次循环重置,上一条消息的残留字段会污染下一条解析结果
以下是可正常运行的实现代码:
package main import ( "encoding/json" "fmt" "io" "os" "google.golang.org/protobuf/proto" "parseRawDHCP/pb" ) func main() { // 以流方式读取文件,避免大文件占满内存 f, err := os.Open("替换为你的Kafka转储文件路径") if err != nil { fmt.Printf("打开文件失败: %v\n", err) return } defer f.Close() buf := make([]byte, 1024*1024) // 1M缓冲区,可根据单条消息最大长度调整 offset := 0 for { // 先读取varint格式的消息长度 l, n := proto.DecodeVarint(buf[offset:]) if n == 0 { // 缓冲区数据不足,读取更多内容 nRead, err := f.Read(buf[offset:]) if err == io.EOF { // 读取完毕正常退出 return } if err != nil { fmt.Printf("读取文件失败: %v\n", err) return } offset += nRead continue } offset += n // 检查缓冲区是否有足够的消息内容 if offset+int(l) > len(buf) { // 剩余内容移到缓冲区头部,再读取新内容 copy(buf, buf[:offset]) nRead, err := f.Read(buf[offset:]) if err != nil { fmt.Printf("读取文件失败: %v\n", err) return } offset += nRead if offset < int(l) { fmt.Printf("消息长度超出缓冲区最大限制,单条消息长度: %d\n", l) return } } // 解析单条消息 var msg pb.Msg err := proto.Unmarshal(buf[offset:offset+int(l)], &msg) if err != nil { fmt.Printf("解析消息失败: %v\n", err) // 跳过损坏的消息继续解析后续内容 offset += int(l) continue } offset += int(l) // 转JSON输出 jsonBytes, err := json.MarshalIndent(&msg, "", " ") if err != nil { fmt.Printf("JSON序列化失败: %v\n", err) continue } fmt.Printf("解析到消息:\n%s\n", jsonBytes) } }
补充说明:
- 如果你的转储文件的长度前缀不是varint格式,而是固定大小的uint32/uint64,只需要把
proto.DecodeVarint替换为对应大小的字节序解析即可 - 废弃的
github.com/golang/protobuf/proto库的Buffer也支持DecodeVarint方法,你如果不想更换依赖可以自行调整对应接口调用 - 代码中加入了损坏消息跳过逻辑,不会因为某一条记录异常导致整个解析流程中断
内容的提问来源于stack exchange,提问作者Igor
相关产品推荐
相关产品推荐

