扇出是将一个输入源数据分发给多个goroutine并行处理,需用闭包封装每个goroutine的读取逻辑以避免竞争和漏数据;扇入是合并多个输出channel,必须用闭包为每个channel绑定专属goroutine确保不丢数据。

什么是扇出(fan-out)?Go里怎么用闭包做并发分发
扇出就是把一个输入源的数据,分发给多个 goroutine 并行处理。关键不是“起多个 goroutine”,而是让每个 goroutine 独立消费同一份输入流,且不互相阻塞或漏数据。
常见错误是直接把同一个 chan int 传给多个函数——这会导致只有一个 goroutine 能读到值,其余阻塞在 <-ch 上,根本没并发。
正确做法是用闭包封装每个 goroutine 的读取逻辑,让每个实例持有自己的读取循环:
func fanOut(in <-chan int, workers int) []<-chan int {
outs := make([]<-chan int, workers)
for i := 0; i < workers; i++ {
ch := make(chan int, 10)
outs[i] = ch
go func(out chan<- int) {
for v := range in { // 注意:这里 in 是共享的,但 range 本身是安全的
out <- v
}
close(out)
}(ch)
}
return outs
}⚠️容易踩的坑:
- 不要用 for range in 在多个 goroutine 里直接读——range 会尝试读完所有值,但 channel 只能被一个 goroutine 消费完
- 闭包捕获变量时别写 go func() { ... }(ch) 而漏掉参数传入,否则所有 goroutine 共享最后一个 ch 的值
- 缓冲区大小设太小(比如 make(chan int, 1))会让扇出侧频繁阻塞,拖慢整体吞吐
扇入(fan-in)为什么必须用闭包?单纯 select 会丢数据
扇入是把多个输出 channel 合并成一个。如果不用闭包隔离每个输入源的读取逻辑,select 在多路复用时会随机选择就绪的 case,导致某个 channel 被反复读、另一个长期饿死——这不是合并,是竞态漏读。
立即学习“go语言免费学习笔记(深入)”;
闭包的作用是为每个输入 channel 绑定专属的 goroutine,确保它持续读取直到关闭:
func fanIn(chs ...<-chan int) <-chan int {
out := make(chan int)
for _, ch := range chs {
go func(c <-chan int) {
for v := range c {
out <- v
}
}(ch) // 必须传参,避免闭包变量捕获问题
}
go func() {
for _, ch := range chs {
<-ch // 等待所有输入 channel 关闭(可选,取决于是否需要精确关闭时机)
}
close(out)
}()
return out
}说明:
- 每个 go func(c <-chan int) 都独占一个 goroutine,不会相互干扰
- out 不能带缓冲(除非你明确要背压),否则可能掩盖下游消费慢的问题
- 如果输入 channel 数量动态变化,别用固定长度的 chs ...<-chan int,改用切片传参 + len(chs) 控制循环
闭包里传 channel 还是传值?看场景选参数类型
闭包参数类型决定并发行为边界。传 <-chan int 表示只读;传 chan<- int 表示只写;传 chan int 表示可读可写——但后者破坏了类型安全,容易引发 panic。
典型组合:
- 扇出侧闭包接收 <-chan int(输入源)和 chan<- int(输出目标)
- 扇入侧闭包只接收 <-chan int(每个输入通道)
- 如果闭包内要关闭 channel,必须传 chan int(且仅限该 goroutine 创建的 channel)
性能影响:
- 传 channel 指针开销极小,远小于复制大量数据
- 但传错方向(比如把 chan<- int 当 <-chan int 用)会在编译时报错,反而是好事
真实场景下,扇入扇出要配合 context 控制生命周期
生产代码里,没人等所有 goroutine 自然结束。比如一个扇出的 worker 因网络超时卡住,整个扇入 channel 就永远发不完。
必须用 context.Context 注入取消信号,且闭包里要监听 ctx.Done():
func worker(ctx context.Context, in <-chan int, out chan<- int) {
for {
select {
case v, ok := <-in:
if !ok {
return
}
result := heavyWork(v)
select {
case out <- result:
case <-ctx.Done():
return
}
case <-ctx.Done():
return
}
}
}注意点:
- ctx.Done() 要在每层 select 里显式检查,不能只在最外层判断
- 不要用 time.Sleep 模拟耗时操作——它不响应 cancel,会绕过 context 控制
- 扇入侧的 goroutine 也要在 for v := range c 前加 select { case <-ctx.Done(): return },否则可能泄漏
复杂点在于,扇出/扇入嵌套越深,context 传递和 cancel 时机越难对齐。最容易被忽略的是:某个中间 channel 关闭后,上游还在往里塞数据,而下游已退出——这时得靠 buffer 或 select default 分流,而不是指望闭包自动处理。


















