fanin需为每个输入channel启独立goroutine配合for range转发,并用sync.waitgroup等待全部完成后再关闭输出channel,避免漏数据或死锁。

fanOut 函数为什么不能直接用 for-range 关闭 results channel
因为 fanOut 中的多个 worker goroutine 是并发写入同一个 results channel 的,如果在主 goroutine 里一收到第一个 job 就关闭 channel,后续 worker 写入会 panic:「send on closed channel」。必须等所有 worker 全部退出后才能安全关闭。
正确做法是用 sync.WaitGroup 计数,每个 worker 执行完 defer wg.Done(),主 goroutine 调用 wg.Wait() 后再 close(results)。
- 别在
fanOut函数末尾直接close(results)—— 这是常见错误 - 如果
results是无缓冲 channel,且没有 goroutine 在接收,worker 第一次写就会阻塞,导致整个池卡死 - 建议把
results设为带缓冲 channel(比如make(chan int, numJobs)),避免 worker 因接收端未就绪而挂住
fanIn 如何合并多个结果 channel 而不漏数据
fanIn 的核心陷阱是:直接用 for range 遍历每个输入 channel,但一旦某个 channel 关闭,range 退出,其他还没读完的 channel 就被跳过了。
标准解法是为每个输入 channel 启一个 goroutine,统一往一个输出 channel 写,再用 sync.WaitGroup 等全部写完后关闭输出 channel:
func fanIn(inputs ...
- 参数用
... 支持任意数量输入 channel,比固定两个更通用 - 不要用
select轮询多个 channel —— 它只取第一个就绪的,无法保证公平或收全 - 如果某个输入 channel 永不关闭(比如持续推送日志),
fanIn就永远不会结束;此时需配合context.Context做超时或取消
pipeline 阶段间要不要加缓冲 channel
加不加取决于阶段处理耗时是否稳定、下游消费速度是否可预期。不加(无缓冲)意味着强同步:上游必须等下游接收后才能发下一个;加了(有缓冲)则能缓解短时波动,但可能掩盖背压问题。
- IO 密集型阶段(如 HTTP 请求)建议用带缓冲 channel,比如
make(chan *http.Response, 10) - CPU 密集型阶段(如图像缩放)慎用大缓冲,否则内存占用飙升,且延迟不可控
- 若某阶段可能 panic 或提前退出,无缓冲 channel 会让上游立即感知阻塞,更容易定位故障点
- 流水线启动时,所有阶段都应显式起 goroutine(
go stage(in)),否则会串行执行,失去 pipeline 意义
context.WithTimeout 怎么嵌入 fanOut/fanIn 流程
单纯在 fanOut 外层套 context.WithTimeout 没用 —— worker 内部的阻塞操作(如 http.Get)不会自动响应 cancel。必须把 ctx 传进每个 worker,并在关键 IO 调用中使用。
- 修改
worker签名:添加ctx context.Context参数,并用ctxhttp.Client替代默认 client -
fanOut启动 worker 时,用ctx派生子 context:childCtx, _ := context.WithCancel(ctx),并在wg.Wait()后调用cancel() -
fanIn的 goroutine 里也要监听ctx.Done(),防止在 channel 关闭前被卡死 - 注意:不要在
fanIn的for n := range c循环里直接 select ctx —— 会中断正常读取;应在每次读之前加select { case
扇出扇入不是堆 goroutine 就完事,真正难的是 channel 生命周期管理 —— 关闭时机、缓冲大小、context 传递路径,这三处出错,轻则丢数据,重则死锁或 panic。
golang免费学习笔记(深入):立即使用
在学习笔记中,你将探索golang的核心概念和高级技巧!











