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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.02 10:27:03