正确关闭通道需用 sync.WaitGroup 同步所有生产者,待其全部完成后再 close(ch),避免 panic 或数据截断;不可由任意生产者自行 close 或主 goroutine 过早 close。
直接用
实现生产者消费者模式本身不难,但线上跑几分钟就卡死、panic 或丢任务,八成是没处理好关闭时机、缓冲策略或 goroutine 协调。
为什么
不能由任意生产者随便调用
向已关闭的
发送数据会触发
。多个生产者共用一个
时,谁关、何时关、是否都关——没有同步机制就会竞态。
错误做法:每个生产者末尾都写
→ 第二个生产者必 panic
错误做法:主 goroutine 在启动生产者后立刻
→ 生产者还没发完,数据被截断
正确做法:用
跟踪所有生产者,再起一个 goroutine 等它们全部
后才
示例关键片段:
带缓冲
的大小不是越大越好
缓冲区本质是内存暂存区,设太大既掩盖背压问题,又可能引发 OOM;设太小则退化为同步阻塞,吞吐上不去。
(无缓冲):收发必须同时就绪,适合强顺序控制,但一端卡住就全链路阻塞
:看似“队列长度为1”,实则和无缓冲行为几乎一致,消费者卡住时第二个生产者仍阻塞
:合理起点,适用于每秒几十到几百任务、单次处理约 100ms 的场景
监控积压用
,但它只是快照,别拿它做调度决策(比如“满一半就降速”不可靠)
消费者怎么安全退出而不漏数据
看似简洁,但只在 channel 被明确关闭后自动退出;如果生产者 panic 退出没来得及
,消费者就永远卡在
里。
立即学习
“
go语言免费学习笔记(深入)
”;
基础安全写法:用
检测关闭
生产级写法:加
和
,支持超时、信号中断、主动取消
示例结构:
注意:不要在
里对同一个
同时做
多消费者共享一个
会不会重复消费
不会。Go 的
是**公平分发**的:多个 goroutine 同时从同一个
读,每次只有一个能拿到数据,底层有 runtime 调度保证,无需额外加锁或标记。
重复消费只发生在「业务层没确认」场景,比如消费者处理完没发 ACK、失败后没重试机制
漏消费常见于:生产者提前
、消费者未用
或
判断就退出、
中漏掉
分支
若需严格“至少一次”语义(如金融任务),就得自己加结果通道、重试计数、死信队列——
本身不提供这些
真正难的从来不是“怎么写通”,而是“怎么在流量突增、消费者卡顿、生产者崩溃时仍不丢不重不 panic”。缓冲大小、关闭时机、退出信号这三处,错一个,整条流水线就不可靠。
chanclose(ch)chanpanic: send on closed channelchanclose(ch)close(ch)sync.WaitGroupDone()close(ch)var producerWg sync.WaitGroup
for i := 0; i < 3; i++ {
producerWg.Add(1)
go func() {
defer producerWg.Done()
for j := 0; j < 5; j++ {
ch <- j
}
}()
}
go func() {
producerWg.Wait()
close(ch) // 唯一合法的关闭位置
}()chanmake(chan int, 0)make(chan int, 1)make(chan int, 100)len(ch)for range chcloserangefor v, ok := ,靠 ok == falsecontext.Contextdone chan struct{}for {
select {
case v, ok := <-ch:
if !ok {
return // channel closed
}
process(v)
case <-ctx.Done():
return // context cancelled
case <-time.After(30 * time.Second):
log.Println("worker timeout")
return
}
}selectchan 和 ch ,容易死锁chanchanchancloserangeokselectok == falsechan