Fan-Out(扇出)是Go并发编程里最实用的模式之一:一个生产者goroutine负责产生数据,多个消费者goroutine并行消费同一个channel,从而把单线程处理不过来的负载摊到多个执行流上。这篇文章从channel的广播特性讲起,手写一个完整的生产者多消费者框架,再讨论负载均衡、优雅退出和goroutine泄漏这些容易踩坑的细节。

Fan-Out的核心原理:一个channel,多个读者
Go的channel有一个天然特性:同一个channel可以被多个goroutine同时读取,每条消息只会被其中一个goroutine取走,不会重复消费。这个语义由运行时调度器保证,相当于内置了一把锁。所以实现扇出最朴素的方式,就是让生产者往一个channel里写数据,然后启动N个goroutine用for range循环读同一个channel。
当生产者关闭channel后,所有阻塞在读取上的消费者会自动退出循环,这是Go设计里非常优雅的一点:for v := range ch会在channel关闭且缓冲取空后结束。也就是说,我们不需要手动通知每个消费者结束,只要在数据生产完毕后调用close(ch),整条流水线就会自然收尾。这一点和某些语言需要显式发送毒丸消息来结束消费者的做法完全不同。
需要注意的是,扇出的前提是消费者之间互相独立,处理顺序不做要求。如果下游必须按生产顺序消费结果,就不能简单地把任务乱序分发给多个消费者,而要在收集端按序号重新组装,这就演变成了另一种模式,这里不展开。
完整实现:生产者、消费者池与结果汇总
下面给出一个完整可运行的例子,模拟一个任务生产者和四个消费者并行处理,并用sync.WaitGroup等待所有消费者完成后汇总结果。这个骨架可以直接套到日志采集、消息批处理、爬虫任务分发等场景。
package main
import (
"fmt"
"sync"
"time"
)
// producer 生产任务,写入完成后关闭 channel
func producer(count int) <-chan int {
out := make(chan int, 10) // 带缓冲,减少生产者阻塞
go func() {
defer close(out)
for i := 1; i <= count; i++ {
fmt.Printf("生产任务 %d\n", i)
out <- i
}
}()
return out
}
// worker 消费者:从 channel 读取任务并处理
func worker(id int, jobs <-chan int, results chan<- string, wg *sync.WaitGroup) {
defer wg.Done()
for job := range jobs { // channel 关闭后自动退出
time.Sleep(100 * time.Millisecond) // 模拟耗时处理
results <- fmt.Sprintf("worker-%d 完成任务 %d", id, job)
}
}
func main() {
jobs := producer(20)
results := make(chan string, 20)
var wg sync.WaitGroup
workerNum := 4
for i := 1; i <= workerNum; i++ {
wg.Add(1)
go worker(i, jobs, results, &wg)
}
// 单独一个 goroutine 负责:等消费者全部结束后关闭 results
go func() {
wg.Wait()
close(results)
}()
// 主 goroutine 收集所有结果
for r := range results {
fmt.Println(r)
}
fmt.Println("全部任务处理完成")
}
这段代码里有几个关键点值得展开。第一,jobs和results都带了缓冲,容量设成任务总量或一个合理值,可以避免生产者和消费者互相卡住;如果任务量不可预估,缓冲给一个中等值即可,靠调度器削峰。第二,关闭results的动作放在了一个单独的goroutine里,先wg.Wait()再close,这个顺序不能反:必须保证没有任何消费者还会往results里写数据之后,才能关闭它,否则会panic(向已关闭的channel发送数据)。
第三,主goroutine用for range results收集结果,配合上面那个自动关闭的goroutine,形成了一个完整闭环。主goroutine不需要自己维护计数器去判断什么时候结束,channel的关闭语义替我们做了这件事。这种谁负责关闭channel的约定是Go并发编程的重要心智模型:通常只有一个发送方时由发送方关闭,多个发送方时由协调方在所有发送方退出后统一关闭。
常见坑点与进阶优化
坑一:goroutine泄漏。如果消费者处理到一半出现return或panic,而生产者还在往一个没人读的channel里写,且缓冲已满,生产者会永久阻塞。解决办法是给每个worker传入context.Context,内部逻辑感知到取消信号就退出,同时生产者也在select中监听ctx的Done,及时停止生产。下面是改造后的worker骨架:
func worker(ctx context.Context, jobs <-chan int) {
for {
select {
case <-ctx.Done():
return // 被取消,立即退出,避免泄漏
case job, ok := <-jobs:
if !ok {
return // channel 已关闭且取空
}
process(job)
}
}
}
坑二:任务粒度和消费者数量不匹配。消费者数量不是越多越好。CPU密集型任务,worker数量设为runtime.NumCPU()就够,开多了只会增加调度开销;IO密集型任务(比如HTTP请求)可以开到几十上百个,但要用带缓冲的channel或信号量限制在途任务总量,防止内存被压爆。可以参考这个经验值:worker数量等于CPU核数乘以一个因子,因子取决于任务阻塞在IO上的时间占比。
坑三:用Fan-In做二次聚合。多个消费者各自往一个共享的results channel里写,本质上已经是扇入的一半。如果把每个消费者的输出channel分开,再用一个聚合goroutine统一读取,结构会更清晰,也方便统计每个消费者的吞吐量。还有一点提醒:永远不要在消费者里关闭任务channel,关闭权归生产者所有,这是channel所有权原则的核心,遵守它能避免绝大多数send on closed channel的panic。
总结一下,Fan-Out模式的骨架就三步:生产者写任务并关闭channel、N个消费者range消费、协调方等所有消费者退出后关闭结果channel。把这个结构吃透,再叠加context取消、动态worker池、限流等能力,就能应对绝大多数Go并发场景了。