用 sync.WaitGroup 控制并发批次最可靠:每批前 wg.Add(1),执行完 defer wg.Done(),主 goroutine 调 wg.Wait() 等本批完成再启下一批,避免资源失控和任务静默丢失。

用 sync.WaitGroup 控制并发批次,别用 for range 直接塞 goroutine
批处理最常崩在资源失控:一次性启几千个 goroutine,内存爆、调度卡、下游扛不住。核心不是“怎么分批”,而是“怎么控并发”。sync.WaitGroup 是最轻量可靠的协调方式,比 channel 控制更直接,也比第三方库少一层抽象风险。
常见错误是把切片按 size 拆好后,对每个子切片起一个 goroutine,却不等它们结束就继续——结果主流程早退了,子任务还在后台飘着。
- 每批启动前调
wg.Add(1),执行完 deferwg.Done() - 主 goroutine 调
wg.Wait()等本批全部完成,再拆下一批 - 别在 goroutine 里 recover panic 后默默吞掉错误;至少 log 并计数,否则失败批次会静默丢失
for i := 0; i < len(data); i += batchSize {
batch := data[i:min(i+batchSize, len(data))]
wg.Add(1)
go func(b []Item) {
defer wg.Done()
processBatch(b) // 实际处理逻辑
}(batch)
}
wg.Wait() // 这行必须有,且在 for 外
bufio.Scanner 读大文件时默认 64KB 缓存不够,得调 Buffer
ETL 经常要读 GB 级日志或 CSV,bufio.Scanner 默认缓存只够单行不超过 64KB。一旦某行超长(比如带大 JSON 字段的 log),直接报 scanner: token too long,整个流程中断。
这不是 bug,是设计取舍:它优先保证安全,不让你无意中加载几 MB 单行进内存。但 ETL 场景下,你得主动接管缓冲区大小。
立即学习“go语言免费学习笔记(深入)”;
- 用
scanner.Buffer(make([]byte, 1024*1024), 1024*1024*4)把初始和最大缓存都设到 4MB - 如果数据源行宽波动极大,建议改用
bufio.Reader.ReadLine()+ 手动拼接,更可控 - 注意:增大缓存不解决内存泄漏——若某行真有 100MB,还是得预判格式并切分,不能全靠调 Buffer
写入 MySQL 批量插入别只用 Exec 拼 SQL,优先走 Prepare + Exec
直拼 SQL 插入百条数据看似快,实则埋雷:SQL 注入风险、参数长度超限、MySQL 的 max_allowed_packet 容易触发、服务端解析压力大。真正稳的批量写法是预编译语句复用。
Go 的 database/sql 对 Prepare 支持良好,但很多人忽略两点:连接复用和参数绑定方式。
- Prepare 在单个
*sql.Conn上调一次即可,反复Exec同一 stmt,别每次重 Prepare - 传参用
[]interface{}切片展开,别手写stmt.Exec(v1, v2, v3)—— 动态批大小下后者难维护 - 如果用
sqlx,db.MustBindStruct或db.Select会自动做字段映射,但批量Insert仍需自己组织VALUES (?, ?, ?)模板
stmt, _ := db.Prepare("INSERT INTO logs (ts, msg, level) VALUES (?, ?, ?)")
for _, b := range batches {
var args []interface{}
for _, item := range b {
args = append(args, item.Timestamp, item.Msg, item.Level)
}
// 注意:args 需按 batch.size * 3 展开
stmt.Exec(args...)
}分批处理中「中间状态」必须持久化,别信内存里的 map 或 slice
ETL 流程跑几小时,中间出错重试时,最头疼的是“从哪续”。有人把已处理 ID 存本地 map,重启就丢;有人靠数据库自增 ID 做断点,但遇到跳号、删数据就错乱。
真正能落地的断点续传,只有一条路:把当前批次的最小/最大游标(如时间戳、offset、ID 范围)写进外部存储,且写操作和本批处理在同一个事务里提交(或至少强顺序 fsync)。
- 推荐用 PostgreSQL 的
INSERT ... ON CONFLICT或 MySQL 的REPLACE INTO记录 checkpoint 表,带唯一约束防重复 - 不要在 batch 处理函数里直接更新 checkpoint;应在 batch 成功写入目标库后,再原子更新 checkpoint
- 如果用消息队列(如 Kafka),offset 提交必须在业务处理成功后,且用
CommitSync,别用异步提交
复杂点从来不在分批逻辑本身,而在“哪批算成功”——这个判断边界一旦模糊,重试就会双写或漏写。盯住 checkpoint 的写入时机和存储可靠性,比优化 for 循环重要十倍。


















