微服务批量处理需解耦“拉、算、存”三阶段并每阶段可控可退可测;禁用串行Exec和OFFSET分页,改用游标分页、事务批量写入及channel流水线;流式解析复用结构体或RawMessage;worker数按瓶颈类型配置,错误须精准捕获与记录。

微服务里处理大批量数据,不是比谁 Goroutine 开得多,而是比谁先卡死在数据库、内存或上下文超时上。核心原则就一条:所有批量操作必须解耦“拉数据”“算数据”“存数据”三阶段,且每阶段可控、可退、可测。
为什么 database/sql 的 Exec 循环在微服务里必崩
微服务默认带 context 超时(常为 30s),而 db.Exec 每条都走独立事务 + 网络 round-trip。万级数据下,哪怕单条 20ms,5000 条就是 100s —— 还没到 DB 就被 context.DeadlineExceeded 中断。更糟的是连接池迅速耗尽,拖垮同实例其他接口。
- 别用
gorm.CreateInBatches顶包:它只是把循环包了一层,底层仍是串行Exec,对 PostgreSQL 比原生pgx.CopyFrom慢 5–10 倍 - MySQL 必须手动拼
INSERT INTO t VALUES (),(),(),但字段含TEXT或索引 > 3 个时,单批别超 200 行,否则触发max_allowed_packet - 所有批量写入前,必须显式
tx := db.Begin(),最后tx.Commit();出错时只对driver.ErrBadConn和net.OpError重试,主键冲突类错误直接记录后跳过
游标分页 + channel 解耦是唯一可行的拉取模式
微服务不能接受 OFFSET 查询的线性延迟恶化。当用户请求“导出全部订单”,后端若用 OFFSET 100000 LIMIT 100,第 1001 页开始每查一次就多扫 10 万行 —— 数据库 CPU 拉满,Go 层还在等 rows.Next() 返回。
- 改用主键游标:
WHERE id > ? ORDER BY id LIMIT ?,每次把上一批最后id当作下一页起点,确保id字段有索引 - 若主键是 UUID,必须建复合索引
(created_at, id),查询条件写成WHERE (created_at, id) > (?, ?) - 拉取和处理必须分离:一个 goroutine 负责
QueryContext并发往inCh chan *Order发数据,N 个 worker 从inCh消费并处理,结果发到outCh;inCh关闭后,主 goroutine 收完outCh再退出 - channel 缓冲区设为
make(chan *Order, 100),太大吃内存,太小导致生产者阻塞在inCh <- order
流式解析 JSONL/CSV 时,别让 Unmarshal 成 GC 灾难
微服务内存敏感,百万行 JSONL 若每行都 json.Unmarshal(data, &v),会触发高频堆分配,pprof 显示 runtime.mallocgc 占比飙升,而 CPU 利用率不足 40%。
Go 配置库,使用 spf13/viper — 分层优先级(flag > env >file > KV > default),提供 BindPFlag/BindPFlags、SetEnvPrefix + SetEnvKeyReplace 等功能。
立即学习“go语言免费学习笔记(深入)”;
- 复用结构体指针:
var record Order; json.Unmarshal(data, &record),但需注意字段值残留(如上一行的record.Name在本行未赋值时仍存在) - 简单结构优先用
json.RawMessage:“先提 key,真要用再解”,例如type Order { ID int; Payload json.RawMessage } - CSV 场景必须用
encoding/csv.Reader.Read(),它内部复用[]string切片;别用strings.Split,每行都新分配切片 - 读文件别用
os.ReadFile:大文件直接 OOM;改用bufio.Scanner,若单行可能超 64KB,提前调scanner.Buffer(make([]byte, 1024*1024), 0)
并发 worker 数量不是越多越好,而是要匹配瓶颈类型
盲目起 100 个 goroutine 处理日志,结果全卡在同一个数据库连接池或 Redis 限流上,调度器反成累赘。
- CPU 密集型(如加签、聚合统计):worker 数 =
runtime.NumCPU(),再多只是增加切换开销 - I/O 密集型(如 HTTP 请求、DB 查询):worker 数 = 连接池大小 × 1.5,例如 pgx 连接池为 20,则 worker 设 30
- 别在
for scanner.Scan()循环里做同步 DB 写入:整个流水线会因单次慢查询卡住;必须把解析后的数据发到 channel,由独立 worker 批量处理 - 每个 worker 处理一批(如 100 条),而非单条 —— 减少 channel 通信频次,实测比单条模式吞吐高 3–5 倍
真正容易被忽略的,是错误传播路径:游标分页中某批 DB 查询失败,若不记录最后成功 id 和失败偏移,重试时就会漏数据或重复;批量入库时某行格式错误,若不单独捕获并落盘原始数据,整批回滚后根本无法定位问题源头。

















