并发模式
本教程共 80 篇 · 第 60 篇 · 更新于 2026-07-27 · 约 9 分钟阅读
60. 并发模式
本节目标:掌握 worker pool、fan-in/fan-out、pipeline、生成器四种经典并发模式,学会用 channel 组织并发流程。
学了 goroutine 和 channel,怎么把它们组合起来解决实际问题?这就是并发模式。这一节讲四个最常用的。
生成器模式
生成器就是用一个 goroutine 往 channel 里生产数据,外部通过 channel 消费:
func generate(nums ...int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for _, n := range nums {
out <- n
}
}()
return out
}
func main() {
for v := range generate(1, 2, 3, 4, 5) {
fmt.Println(v)
}
}
函数返回一个只读 channel,内部 goroutine 负责填数据并关闭。调用方用 range 消费,干净利落。
这个模式的好处是「把生产过程封装起来」,调用方不用关心数据怎么来的。
Pipeline 模式
pipeline(流水线)把任务拆成多个阶段,每个阶段用一个 goroutine 处理,阶段之间用 channel 连接:
// 阶段1:生成数字
func gen(nums ...int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for _, n := range nums {
out <- n
}
}()
return out
}
// 阶段2:平方
func square(in <-chan int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for n := range in {
out <- n * n
}
}()
return out
}
// 阶段3:打印
func main() {
for v := range square(gen(1, 2, 3, 4)) {
fmt.Println(v) // 1 4 9 16
}
}
每个阶段读上一个阶段的 channel,处理后写到下一个 channel。数据像流水一样经过每一站。
Tippipeline 的关键:每个阶段函数的输入和输出都是 channel。这样就能随意组合,像积木一样拼起来。
Worker Pool 模式
worker pool 是一组固定数量的 worker goroutine,从一个共享 channel 取任务执行。好处是控制并发数量,避免开太多 goroutine 把资源撑爆:
func worker(id int, jobs <-chan int, results chan<- int) {
for j := range jobs {
fmt.Printf("worker %d 处理任务 %d\n", id, j)
results <- j * j
}
}
func main() {
jobs := make(chan int, 10)
results := make(chan int, 10)
// 启动 3 个 worker
for w := 0; w < 3; w++ {
go worker(w, jobs, results)
}
// 发送 5 个任务
for j := 1; j <= 5; j++ {
jobs <- j
}
close(jobs)
// 收集结果
for r := 0; r < 5; r++ {
fmt.Println("结果:", <-results)
}
}
jobs channel workers results channel
[1,2,3,4,5] --> worker0 --> [1,4,9,16,25]
--> worker1
--> worker2
worker 数量固定为 3,无论任务多少,同时只有 3 个在跑。
Noteworker pool 适合任务量大但每个任务都不重的场景。如果任务之间有依赖或需要限速,可以配合
context做取消,用time.Ticker做限流。
Fan-out / Fan-in 模式
fan-out 是把一个 channel 的数据分发给多个 goroutine 并行处理;fan-in 是把多个 goroutine 的结果汇总到一个 channel:
// fan-out:把输入分给多个 worker
func split(in <-chan int, n int) []<-chan int {
outs := make([]<-chan int, n)
for i := 0; i < n; i++ {
ch := make(chan int)
outs[i] = ch // 双向自动转只读,存进切片
go func(out chan<- int) {
defer close(out)
for v := range in {
out <- v
}
}(ch) // 传双向 ch,函数参数收成只写
}
return outs
}
// fan-in:把多个 channel 合并成一个
func merge(chs ...<-chan int) <-chan int {
out := make(chan int)
var wg sync.WaitGroup
wg.Add(len(chs))
for _, ch := range chs {
go func(c <-chan int) {
defer wg.Done()
for v := range c {
out <- v
}
}(ch)
}
go func() {
wg.Wait()
close(out)
}()
return out
}
func main() {
in := gen(1, 2, 3, 4, 5, 6)
outs := split(in, 3) // fan-out 到 3 个 worker
for v := range merge(outs...) { // fan-in 合并结果
fmt.Println(v)
}
}
fan-out 提升并行度,fan-in 收拢结果。merge 里用 sync.WaitGroup 等所有 worker 结束再 close,这个技巧很常用。
超时与取消
实际程序里,你得能随时叫停。配合 select 和 context 做超时取消:
func main() {
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
defer cancel()
ch := make(chan int)
go func() {
// 模拟慢任务
time.Sleep(3 * time.Second)
ch <- 42
}()
select {
case v := <-ch:
fmt.Println("收到:", v)
case <-ctx.Done():
fmt.Println("超时取消:", ctx.Err())
}
}
context 下一部分会详细讲,这里先知道它能配合 select 做超时就行。
模式选择的思路
| 模式 | 适用场景 |
|---|---|
| 生成器 | 封装数据生产过程 |
| Pipeline | 任务有多个处理阶段,可串行流水化 |
| Worker Pool | 任务量大,需要控制并发数 |
| Fan-out/Fan-in | 单个阶段想并行加速 |
Warning别为了用模式而用模式。简单的串行调用能搞定的事,硬套并发模式反而增加复杂度。channel 和 goroutine 也有开销,数据量小的时候串行更快。
小结
- 生成器:goroutine 往 channel 生产数据,封装生产逻辑
- Pipeline:多阶段串成流水线,阶段间用 channel 连接
- Worker Pool:固定数量 worker 从任务 channel 取活干
- Fan-out/Fan-in:分散并行处理,再汇总结果
- 用
sync.WaitGroup等待多个 goroutine 完成后关闭 channel
下一节讲数据竞争和竞态检测器,这是写并发程序绕不过的坑。