流式读取的核心目标是在不把整个输入一次性读入内存的前提下,按某种边界把数据切成一段段可处理的内容。Go 标准库的 bufio.Scanner 已经把这个流程包装得很好,但默认只支持按行、按单词等简单规则切分。如果日志或协议里使用 |||、END 这样的多字节字符串作为记录分隔符,就需要自己扩展切分逻辑。先看 Scanner 的 SplitFunc 机制。

理解 bufio.Scanner 的 SplitFunc 接口
bufio.Scanner 把输入拆成一个一个令牌的过程,本质上是一个循环:每次从底层 io.Reader 读入一批字节,交给一个分割函数判断从哪里切开,然后返回一个记录。这个分割函数的类型是 bufio.SplitFunc,签名如下:func(data []byte, atEOF bool) (advance int, token []byte, err error)。data 是当前缓冲区中尚未被消费的数据;atEOF 表示底层读取器是否已经到达 io.EOF;advance 用来告诉 Scanner 本次可以消费多少字节,token 是切分出来的记录,err 则用于终止扫描。
默认的 ScanLines 会把 data 中的第一个换行符当成边界,命中后返回换行符之前的内容,并消费掉换行符。这个逻辑只认一个或两个字节的换行符,完全无法处理像 ||| 这样的多字节分隔符。于是最简单的思路就是照葫芦画瓢:自己实现一个 SplitFunc,在 data 中查找指定的字节序列,找到就切分,找不到就告诉 Scanner 等待更多数据。
这里有一个容易忽略的返回约定:当函数返回 0, nil, nil 时,Scanner 会继续读取更多数据并再次调用分割函数,且之前的数据仍然保留在缓冲区中。当返回的 advance 大于 0 时,Scanner 会把对应字节标记为已消费;如果 token 为 nil 且 err 为 nil,则本次不会产出记录,而是继续扫描。理解这些返回值,是实现稳定分割函数的基础。
用 bytes.Index 编写多字节分割函数
实现一个按字符串分隔的分割函数并不复杂。下面的 splitByString 函数接收一个分隔符字符串,返回一个闭包形式的 bufio.SplitFunc。闭包中把分隔符转换成 []byte,然后在每次调用时用 bytes.Index 搜索它的位置。
package main
import (
"bufio"
"bytes"
"fmt"
"strings"
)
func splitByString(delim string) bufio.SplitFunc {
d := []byte(delim)
return func(data []byte, atEOF bool) (advance int, token []byte, err error) {
if atEOF && len(data) == 0 {
return 0, nil, nil
}
if i := bytes.Index(data, d); i >= 0 {
return i + len(d), data[:i], nil
}
if atEOF {
return len(data), data, nil
}
return 0, nil, nil
}
}
func main() {
input := "hello|||world|||go"
scanner := bufio.NewScanner(strings.NewReader(input))
scanner.Split(splitByString("|||"))
for scanner.Scan() {
fmt.Printf("token: %q\n", scanner.Text())
}
if err := scanner.Err(); err != nil {
fmt.Println("scan error:", err)
}
}
这个函数的核心逻辑分成四步。第一步处理输入已经结束且缓冲区为空的情况,返回 0, nil, nil 表示没有更多记录。第二步在 data 中查找分隔符,一旦找到 i,就让 advance 等于分隔符起点加分隔符长度,同时把分隔符之前的内容作为 token 返回。这样 Scanner 会跳过整个分隔符,下一次扫描从分隔符之后的第一个字节继续。
如果还没有找到分隔符,但 atEOF 为 true,说明输入已经结束,剩余数据就是最后一条记录,直接返回全部剩余字节即可。否则返回 0, nil, nil,请求 Scanner 读入更多数据。这样无论输入是 hello|||world|||go,还是非常大的文件,都能按 ||| 稳定切分。使用时只需要一行 scanner.Split(splitByString("|||")),其余逻辑跟按行扫描完全一致。
这种方式非常适合记录长度可控、分隔符出现频率较高的场景。但要注意,Scanner 默认单条记录的最大长度是 64KB。如果两条分隔符之间的内容超过这个阈值,就会触发 bufio.ErrTooLong。可以通过 scanner.Buffer 调整缓冲区上限,不过这只是提高上限,并不能改变分割函数在未命中时数据不断堆积的本质。
用 bufio.Reader 维护可控缓冲区
上面的方案依赖 Scanner 内部的缓冲区。当分隔符尚未出现时,分割函数只能返回 0, nil, nil,所有已读入的数据都会一直保留在缓冲区中。对于超大记录,即使调大了 MaxScanTokenSize,仍然可能因为分隔符出现得太晚而占用大量内存。如果想要更细粒度地控制读取过程,并且实现真正的边读边切,可以改用 bufio.Reader 自己维护一个窗口。
下面的 StringDelimReader 结构体封装了一个 bufio.Reader、目标分隔符和一块中间缓冲区。每次调用 Next 时,它先在缓冲区中查找分隔符;如果找到,就切出前面的内容,并保留分隔符之后的剩余部分,等待下一次调用。如果没找到,就继续从底层读取固定大小的块追加到缓冲区。这样不会把整个输入一次性吞进来,而是始终保持一个较小的待处理窗口。
package main
import (
"bufio"
"bytes"
"fmt"
"io"
"strings"
)
type StringDelimReader struct {
reader *bufio.Reader
delim []byte
buffer []byte
maxSize int
eof bool
}
func NewStringDelimReader(r io.Reader, delim string, maxSize int) *StringDelimReader {
return &StringDelimReader{
reader: bufio.NewReader(r),
delim: []byte(delim),
maxSize: maxSize,
}
}
func (s *StringDelimReader) Next() ([]byte, error) {
for {
if i := bytes.Index(s.buffer, s.delim); i >= 0 {
token := s.buffer[:i]
s.buffer = s.buffer[i+len(s.delim):]
return token, nil
}
if s.eof {
if len(s.buffer) > 0 {
token := s.buffer
s.buffer = nil
return token, nil
}
return nil, io.EOF
}
if len(s.buffer) > s.maxSize {
return nil, fmt.Errorf("token exceeds max size %d", s.maxSize)
}
chunk := make([]byte, 4096)
n, err := s.reader.Read(chunk)
if n > 0 {
s.buffer = append(s.buffer, chunk[:n]...)
}
if err == io.EOF {
s.eof = true
} else if err != nil {
return nil, err
}
}
}
func main() {
input := "alpha|||beta|||gamma"
r := NewStringDelimReader(strings.NewReader(input), "|||", 1024*1024)
for {
token, err := r.Next()
if err == io.EOF {
break
}
if err != nil {
fmt.Println("read error:", err)
break
}
fmt.Printf("token: %q\n", token)
}
}
这个结构的优点是明确控制内存上限:maxSize 用来限制缓冲区最大长度,一旦超过就返回错误,避免异常输入导致内存无限膨胀。同时,它不像 Scanner 那样必须等一条记录完全结束时才把它交给调用方,而是可以在循环中逐步处理。不过代价也很明显:我们需要自己管理 eof 标志、剩余缓冲区和错误返回值,代码比直接实现 SplitFunc 更复杂。
另一个需要注意的细节是,这里没有把分隔符后面的数据“放回”读取器,因为 bufio.Reader 一旦调用 Read 消费了字节,就无法直接回退。所以命中的时候必须把分隔符后面的数据继续留在 buffer 中,作为下一次 Next 的起点。这是手动维护缓冲区的核心约束,也是很多实现里容易写错的地方。
边界情况、错误处理与参数调优
多字节分隔符跨缓冲区边界是一个常见的担忧。使用 Scanner 方案时,只要分隔符尚未被完整找到,分割函数就会持续请求更多数据,因此跨内部缓冲边界的情况通常由 Scanner 自动处理。但前提是单条记录长度不能超过配置的最大令牌长度。可以通过 scanner.Buffer(make([]byte, 0, 64*1024), 4*1024*1024) 把上限调整到 4MB 甚至更大,避免处理较大日志记录时触发错误。
空令牌问题也需要根据业务决定是否过滤。如果输入以分隔符开头,例如 |||first|||second,第一次扫描会返回一个空字符串。对于多数协议而言,空记录可能意味着无效数据,可以在循环中加入 if len(scanner.Text()) > 0 进行过滤。但如果业务需要保留空记录的位置,就不能简单跳过,否则后续记录的下标会与原始流不一致。
性能方面,bytes.Index 在 Go 标准库中已经做了充分优化,对于常见的短分隔符和中等长度的缓冲区,直接使用它通常足够快。如果分隔符非常长,或者在极高频调用中出现重复扫描的开销,可以考虑使用 KMP 等预计算前缀表的方式优化查找。但大多数实际应用中,网络或磁盘 I/O 才是主要瓶颈,分割逻辑本身并不需要过度优化。
最后做一个简单选择:如果记录规模在几 MB 以内,并且希望代码尽量少,直接扩展 bufio.Scanner 是最省事的方式;如果数据流中记录可能非常大,或者需要严格控制内存占用,手动用 bufio.Reader 维护缓冲区会更可靠。两种方案本质上都是把一段多字节字符串当作边界标记,再根据边界把字节流逐步转换成可以逐条处理的记录。