fanin 函数必须接收 []interface{} 类型参数,用于合并多个通道的数据流,实现并发控制与数据聚合。

fanIn 函数必须接收 [],不能是 <code>[]chan T
类型声明错一个箭头,就会让调用方误传可写 channel,导致编译通过但运行时 panic:send on closed channel 或 goroutine 泄漏。Go 的 channel 方向是类型系统的一部分, 表示“只读”,确保 fanIn 只消费、不生产;而 <code>chan T 是双向,调用方可能意外往里塞数据,破坏归集逻辑。
- 正确签名:
func fanIn[T any](cs [] - 错误签名:
func fanIn(cs []chan T) chan T—— 这会让使用者误以为可以往每个 cs[i] 写,实际 fanIn 只该读 - 如果输入 channel 来自不同 goroutine(如多个 HTTP 请求),它们各自关闭是安全的;fanIn 不负责关输入,只等它们关完再关自己的 out
每个输入 channel 必须由独立 goroutine 转发,不能用 select 轮询固定列表
用 for { select { case v := 手动展开所有 channel,看似简单,实则不可扩展、易丢数据:一旦某个 ch 关闭,对应 case 会持续就绪,<code>select 可能反复选中它(因 default 缺失或逻辑错误),导致其他 channel 长期饥饿。更糟的是,新增一路输入就得改函数体,违反开闭原则。
- 正确做法:为每个
启一个 goroutine,只做一件事——读完本 channel 全部数据,转发到统一 <code>out - 转发 goroutine 结束后自动退出,不关
out,也不影响其他路 - 主 fanIn goroutine 用
sync.WaitGroup等所有转发 goroutine 完毕,再close(out)
必须用 sync.WaitGroup 等待全部输入结束,不能靠 range 或超时硬切
常见错误是启动转发 goroutine 后,直接在 fanIn 主 goroutine 里 for range out —— 这会阻塞等待第一个值,但此时其他路还没启动完,或者某路卡住没发数据,整个流程就挂死。更危险的是用 time.After 强制关闭 out,结果部分数据还在路上就被截断。
- fanIn 内部必须起一个额外 goroutine 调用
wg.Wait(),然后close(out) -
wg.Add(len(cs))在启动每个转发 goroutine 前调用,defer wg.Done()在转发 goroutine 末尾 - 不要在转发 goroutine 里 close 输入 channel —— 它们应由上游生产者关闭;fanIn 只管“收完即关 out”
缓冲大小设为 len(cs) 就够了,别盲目加大
很多人以为 output channel 缓冲越大越不容易阻塞,其实不然。fanIn 的输出是严格按接收顺序来的,缓冲过大(比如 make(chan int, 1000))会让主 goroutine 误判“数据已收齐”,实际只是被缓住了,下游消费慢时反而掩盖背压,拖垮整体响应。尤其当某路输入特别快、其他路极慢时,大缓冲会让快路数据囤积,内存占用陡增。
- 推荐缓冲大小:
make(chan T, len(cs))—— 最多暂存每路一个值,既防 goroutine 阻塞,又暴露真实延迟 - 如果下游处理耗时波动大,宁可加
context.WithTimeout控制 fanIn 整体生命周期,也不靠 buffer 掩盖问题 - 无缓冲 channel(
make(chan T))在低吞吐场景可行,但高并发下会显著降低吞吐,因为每次写入都得等下游 ready
wg.Done(),或早关了一次 out,整条流水线就卡在半空。golang免费学习笔记(深入):立即使用
在学习笔记中,你将探索golang的核心概念和高级技巧!











