不能直接 for range channel 拉消息,因为其只检测通道关闭而不等待 worker 完成;主 goroutine 退出会导致正在处理的 worker 被强制终止,引发数据丢失或 panic;必须用 sync.WaitGroup 显式管理生命周期,并配合 context 控制长连接退出。

为什么不能直接 for range channel 拉消息就完事
因为 RabbitMQ 或 Redis 消费者拉到的消息是流式的,for range 只管“通道关了没”,不管“worker 是否还在处理”。主 goroutine 一退出,所有正在 time.Sleep、DB 写入、HTTP 调用的 worker 全被杀掉,日志只打印一半,消息看似“消费了”实则丢了。
常见错误现象:Worker 2 processing task 7: data-7 打到一半程序静默退出;或者 panic: send on closed channel 来自某个 worker 想写结果回 channel 时发现已被关闭。
- 必须用
sync.WaitGroup显式计数:每个 worker 启动前wg.Add(1),结束前defer wg.Done() - 生产者(比如 RabbitMQ 的
msgs接收循环)负责单点close(ch),消费者只读,绝不关 - 若消费者是长连接(如监听 Kafka topic),就不能依赖
close,得换context.Context控制退出
batchProcessor 怎么写才不丢数据、不爆内存
批量处理的核心不是“攒够 N 条就发”,而是“要么全成功,要么全重试”,否则 ES 写一半失败、MySQL 插一半回滚,状态就撕裂了。Golang 里最稳的方式是让每个 batch 在单个 goroutine 内串行执行,天然避免锁竞争。
使用场景:RabbitMQ 消费后批量写 Elasticsearch、Redis Stream 消费后批量更新缓存、Kafka 消息聚合后触发风控规则。
立即学习“go语言免费学习笔记(深入)”;
-
BatchSize别硬设成 100 —— 实际要看下游吞吐:ES bulk 推荐 5–20MB / request,换算下来常是 500–2000 条/批;MySQL 批量 INSERT 则建议 ≤ 1000 行 - 超时控制必须用
time.Timer,别用time.After:后者在select循环里会持续创建新 timer,导致 goroutine 和内存泄露 - 每批处理完要手动
Ack(RabbitMQ)或XAck(Redis Stream),不能靠自动确认——否则处理失败时消息已丢
prefetch count 和 Workers 数怎么配才不堆积
RabbitMQ 的 prefetch count 和 Golang 里的 Workers 数量不是一回事,但协同不好就会一边疯狂积压 Ready 消息,一边 Unacked 卡死一堆。
参数差异:prefetch count = 10 表示 RabbitMQ 最多给这个 consumer 发 10 条未确认消息;而 Workers = 4 是你本地开了 4 个 goroutine 并发处理——如果每个 worker 处理慢,这 10 条就全卡在 Unacked 状态,新消息进不来。
- 起步值建议:RabbitMQ
prefetch count设为Workers × 2(比如 4 个 worker 就设 8),留缓冲余量 - Workers 数别盲目堆高:超过
runtime.NumCPU()的 2–4 倍后,调度开销反升,尤其当业务含 DB 查询、HTTP 调用等阻塞操作时 - 监控关键指标:RabbitMQ Management UI 里看
Unacked是否长期 >prefetch count;Golang 进程看goroutines数是否稳定(突增可能有 leak)
channel 缓冲大小设多少才不卡又不 OOM
make(chan Message, N) 的 N 不是并发数,是“等待区”容量。它决定了生产者能一口气塞多少消息进去而不阻塞——但塞进去不代表有人处理,只是暂时存着。
性能影响很直接:设太小(如 N=1),突发流量下生产者频繁阻塞,吞吐断崖下跌;设太大(如 N=10000),内存占用飙升,且掩盖真实瓶颈(你以为是队列满,其实是 worker 卡在 MySQL INSERT 上)。
- 经验值:从
100起步;若日志高频出现len(ch) == cap(ch),说明消费者跟不上,优先优化batchProcessor内部逻辑,而不是扩 buffer - 结构体大小要算:假设
Message平均 2KB,N=1000就占 2MB 内存——百个消费者就是 200MB,得看服务资源 - 允许丢弃?可用
select { case ch 主动背压,比 OOM 强
真正难的不是写对语法,是搞清哪一层在卡:是网络 IO?DB 连接池不够?还是 ES bulk size 超限被拒?得一层层查指标,而不是一上来就加 goroutine。


















