Go 高手(16):并发模式——worker pool、pipeline、fan-in/fan-out
更新时间:2026-09-01。本文是
languages/go/intermediate/高手层第 16 篇,接 Go 工具链。Go 并发编程有一些经典模式,这些模式不是设计模式,而是 goroutine + channel 的常见组合方式。掌握了它们,就能处理绝大多数并发场景。
本文要回答的问题
- worker pool 模式怎么实现?为什么用 channel 做信号量?
- pipeline 模式怎么串联多个 stage?stage 之间怎么传递数据?
- fan-out 和 fan-in 有什么区别?什么时候用?
- 并发模式的超时和取消怎么处理?
一、Worker Pool(工作池)
控制并发数量,避免无限制的 goroutine 消耗资源:
func workerPool(jobs []Job, concurrency int) []Result {
jobCh := make(chan Job, len(jobs))
resultCh := make(chan Result, len(jobs))
var wg sync.WaitGroup
// 启动 worker
for i := 0; i < concurrency; i++ {
wg.Add(1)
go func() {
defer wg.Done()
for job := range jobCh {
resultCh <- process(job)
}
}()
}
// 发送任务,然后关闭
for _, job := range jobs {
jobCh <- job
}
close(jobCh)
// 等待所有 worker 完成
wg.Wait()
close(resultCh)
// 收集结果
var results []Result
for res := range resultCh {
results = append(results, res)
}
return results
}核心思想: 固定数量的 goroutine 从 channel 取任务,处理完放结果。控制 concurrency 参数就能控制并发数。
二、Pipeline(流水线)
多个 stage 串联,每个 stage 用 goroutine 处理,通过 channel 连接:
// stage 1: 生成数字
func gen(nums ...int) <-chan int {
out := make(chan int)
go func() {
for _, n := range nums {
out <- n
}
close(out)
}()
return out
}
// stage 2: 平方
func sq(in <-chan int) <-chan int {
out := make(chan int)
go func() {
for n := range in {
out <- n * n
}
close(out)
}()
return out
}
// stage 3: 过滤
func filter(in <-chan int, threshold int) <-chan int {
out := make(chan int)
go func() {
for n := range in {
if n > threshold {
out <- n
}
}
close(out)
}()
return out
}
// 使用
c := gen(1, 2, 3, 4, 5)
c = sq(c)
c = filter(c, 10)
for v := range c {
fmt.Println(v) // 16, 25
}pipeline 的优点:
- 每个 stage 独立,可以单独测试
- 每个 stage 可以独立并发
- 新增 stage 只需要加一个 goroutine 函数
三、Fan-out / Fan-in
Fan-out:一个 channel 的数据分发到多个 goroutine 处理。 Fan-in:多个 goroutine 的结果合并到一个 channel。
// fan-out:一个 channel 分发到多个 worker
func fanOut(in <-chan int, numWorkers int) []<-chan Result {
channels := make([]<-chan Result, numWorkers)
for i := 0; i < numWorkers; i++ {
ch := make(chan Result)
go func() {
defer close(ch)
for v := range in {
ch <- process(v)
}
}()
channels[i] = ch
}
return channels
}
// fan-in:多个 channel 合并到一个
func fanIn(channels ...<-chan Result) <-chan Result {
out := make(chan Result)
var wg sync.WaitGroup
for _, ch := range channels {
wg.Add(1)
go func(c <-chan Result) {
defer wg.Done()
for v := range c {
out <- v
}
}(ch)
}
go func() {
wg.Wait()
close(out)
}()
return out
}典型场景: 爬虫。一个 URL 列表 fan-out 到多个 worker 并发爬取,结果 fan-in 到一个 channel 集中处理。
四、并发超时控制
func processWithTimeout(items []Item, timeout time.Duration) ([]Result, error) {
ctx, cancel := context.WithTimeout(context.Background(), timeout)
defer cancel()
resultCh := make(chan Result, len(items))
for _, item := range items {
item := item
go func() {
result, err := process(item)
select {
case resultCh <- Result{Item: item, Result: result, Err: err}:
case <-ctx.Done():
// 超时了,不发送了
}
}()
}
// 等所有结果或超时
var results []Result
for i := 0; i < len(items); i++ {
select {
case res := <-resultCh:
results = append(results, res)
case <-ctx.Done():
return results, ctx.Err()
}
}
return results, nil
}五、Pipeline 的取消传递
func gen(ctx context.Context, nums ...int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for _, n := range nums {
select {
case out <- n:
case <-ctx.Done():
return
}
}
}()
return out
}每个 stage 都检查 ctx.Done(),context 取消时整个 pipeline 立即停止,不浪费资源。
六、常见坑对照
| 坑 | 现象 | 对策 |
|---|---|---|
| 忘记 close channel | 接收方 range 永远不退出,死锁 | 发送方确认数据发完后 close |
| pipeline 中某个 stage 卡住 | 整个 pipeline 阻塞 | 每个 stage 都支持 context 取消 |
| channel 发送不匹配接收 | 死锁 | 确保发送和接收数量匹配 |
| 无限 goroutine 创建 | 资源耗尽,OOM | 用 worker pool 控制并发数 |
相关与延伸
下一篇:错误处理模式——哨兵错误、自定义错误、错误包装最佳实践;进阶层更深入的并发原语,见 goroutine 与 channel 深度。
一句话总结
Go 并发模式:worker pool 用固定 goroutine + channel 控制并发数;pipeline 用多个 goroutine 串联,每个 stage 独立并发;fan-out 一个 channel 分发到多个 goroutine,fan-in 多个 channel 合并到一个;所有并发模式都要支持 context 取消,避免资源泄漏;close channel 通知接收方结束,必须由发送方 close。