当数据源体积大、内存敏感或下游处理速度不稳时,应用channel分批流式处理;需设小缓冲(如16)、显式控制批次大小,并在消费者端加recover和超时保护。

什么时候该用 channel 做分批流式处理,而不是一次性读完?
当你面对的数据源体积大(比如 GB 级日志文件、数据库游标结果集)、内存敏感、或下游消费者处理速度不稳时,channel 是 Go 中最自然的流式分批载体。它天然支持协程间解耦、背压传递和异步缓冲——但前提是别把 channel 当成无界队列用。
常见错误是直接 make(chan []T, 0) 或 make(chan []T, 100000):前者阻塞严重,后者等于在内存里缓存整批数据,失去“流式”意义。
- 推荐用带小缓冲的
channel(如make(chan []T, 16)),缓冲大小 ≈ 下游单次处理耗时 × 预估吞吐波动幅度 - 上游生产者必须显式控制每批次大小(例如按行数、字节数或时间窗口切分),不能依赖
range自动拆分 - 若下游处理可能 panic 或超时,需在消费者端加
recover和select超时,否则整个流水线会卡死
bufio.Scanner + channel 分批读文件的典型写法
这是最常被误用的组合:bufio.Scanner 默认按行扫描,但它的 Scan() 是同步阻塞的,不能直接扔进 goroutine 后就不管了。必须手动控制“一批读多少行”,再推到 channel。
func batchLines(filename string, batchSize int) 0 {
out
- 关键点:
scanner.Bytes()返回的是复用的底层缓冲区,下一次Scan()就会覆盖——必须append([]byte(nil), scanner.Bytes()...)或用string(scanner.Text())转成独立副本 - 错误现象:
batch里所有元素最后都变成最后一行内容,就是没做深拷贝 - 如果按字节而非行分批,改用
io.ReadFull或bufio.NewReader的Read更可控
如何让分批处理支持动态批大小和取消信号?
硬编码 batchSize 在真实场景中很快会失效:上游速率突增、下游限流策略变化、或用户主动中断任务。这时需要把批大小和上下文控制权交给调用方。
正确做法是把 context.Context 和参数结构体一起传入,且 channel 类型保持为 (只读),避免消费者意外关闭或写入。
type BatchOptions struct {
BatchSize int
Timeout time.Duration
}
<p>func StreamBatches(ctx context.Context, src io.Reader, opt BatchOptions) 0 {
out </p>
-
context.Context必须在每次循环开始检查,不能只在 goroutine 开头 check 一次 - 不要用
time.After替代ctx.Done(),否则无法响应父 context 的取消 - 如果批大小需运行时调整(如根据 CPU 负载自动缩放),就得换成带回调的工厂函数,而非静态参数
数据库查询结果怎么安全地分批流式消费?
用 database/sql 的 Rows 做流式处理时,最大的陷阱是:你以为 rows.Next() 是流式,其实底层驱动可能已把整张表拉到内存(尤其 MySQL 驱动默认不启用流式游标)。
- PostgreSQL:确保连接串含
pgx驱动 +prefer_simple_protocol=true,并用QueryRow或带rows.Close()的显式循环 - MySQL:必须设置
mysql.ParseTime=true&loc=Local并确认驱动版本 ≥ 1.6,否则sql.Rows会缓冲全部结果 - 通用原则:每批处理完立即调用
rows.Scan(),不要累积[][]interface{},更不要用rows.SliceScan()
一个安全的分批封装示例:
func QueryBatches(ctx context.Context, db *sql.DB, query string, args ...interface{}) <p>真正流式的关键不在 Go 代码,而在驱动是否真正逐行 fetch——这点很容易被忽略,直到 OOM 才发现。</p>











