Appearance
并发模式:WorkerPool 与 Pipeline
前面几篇我们学习了 goroutine、channel、select、sync 包这些「积木」。本篇我们将把它们组合起来,学习 Go 并发编程中最经典的几种模式:Worker Pool、Pipeline、Fan-out/Fan-in、Generator、Future/Promise、Timeout、Cancel。这些模式是前人总结出的并发设计经验,掌握它们能让你在面对并发问题时快速找到合适的结构。
一、Worker Pool 模式
Worker Pool(工作池)模式:启动固定数量的 worker goroutine,从一个共享的任务 channel 取任务执行,结果写到结果 channel。这种模式能控制并发数,避免无限制创建 goroutine。
1. 模式结构
任务生产者 → [任务 channel] → N 个 worker → [结果 channel] → 结果消费者2. 完整实现
go
package main
import (
"fmt"
"sync"
"time"
)
type Task struct {
ID int
Input int
}
type Result struct {
TaskID int
Output int
}
func worker(id int, tasks <-chan Task, results chan<- Result, wg *sync.WaitGroup) {
defer wg.Done()
for task := range tasks {
fmt.Printf("worker %d 处理任务 %d\n", id, task.ID)
// 模拟计算
time.Sleep(200 * time.Millisecond)
results <- Result{
TaskID: task.ID,
Output: task.Input * task.Input, // 平方
}
}
}
func main() {
const numWorkers = 3
const numTasks = 10
tasks := make(chan Task, numTasks)
results := make(chan Result, numTasks)
var wg sync.WaitGroup
// 启动固定数量的 worker
for w := 1; w <= numWorkers; w++ {
wg.Add(1)
go worker(w, tasks, results, &wg)
}
// 投递任务
for i := 1; i <= numTasks; i++ {
tasks <- Task{ID: i, Input: i}
}
close(tasks) // 关闭任务 channel,worker 处理完会自动退出
// 等所有 worker 完成,然后关闭 results
go func() {
wg.Wait()
close(results)
}()
// 收集结果
for r := range results {
fmt.Printf("结果: 任务 %d → %d\n", r.TaskID, r.Output)
}
fmt.Println("全部完成")
}3. Worker Pool 的优点
- 控制并发数:worker 数量固定,不会无限增长。
- 复用 goroutine:一个 worker 处理多个任务,避免频繁创建销毁。
- 削峰:任务 channel 缓冲可以平滑突发的任务高峰。
4. 带优雅退出的 Worker Pool
go
package main
import (
"context"
"fmt"
"sync"
"time"
)
func worker(ctx context.Context, id int, tasks <-chan int, wg *sync.WaitGroup) {
defer wg.Done()
for {
select {
case <-ctx.Done():
fmt.Printf("worker %d 被取消,退出\n", id)
return
case task, ok := <-tasks:
if !ok {
fmt.Printf("worker %d 任务结束,退出\n", id)
return
}
fmt.Printf("worker %d 处理 %d\n", id, task)
time.Sleep(100 * time.Millisecond)
}
}
}
func main() {
ctx, cancel := context.WithCancel(context.Background())
tasks := make(chan int, 10)
var wg sync.WaitGroup
// 启动 3 个 worker
for i := 1; i <= 3; i++ {
wg.Add(1)
go worker(ctx, i, tasks, &wg)
}
// 投递一些任务
for i := 1; i <= 5; i++ {
tasks <- i
}
// 模拟提前取消
time.Sleep(300 * time.Millisecond)
fmt.Println("发起取消")
cancel()
wg.Wait()
fmt.Println("所有 worker 退出")
}二、Pipeline 模式
Pipeline(流水线)模式:把一个复杂处理拆成多个阶段,每个阶段是一个 goroutine,通过 channel 串联。前一个阶段的输出是后一个阶段的输入,像工厂流水线一样。
1. 模式结构
阶段1 → [chan] → 阶段2 → [chan] → 阶段3 → 结果2. 完整实现:数据处理流水线
go
package main
import "fmt"
// 阶段1:生成数据
func generate(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 filterEven(in <-chan int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for n := range in {
if n%2 == 0 {
out <- n
}
}
}()
return out
}
// 阶段4:打印
func printer(in <-chan int) {
for n := range in {
fmt.Println("输出:", n)
}
}
func main() {
// 串联:generate → square → filterEven → printer
nums := generate(1, 2, 3, 4, 5, 6)
squared := square(nums)
even := filterEven(squared)
printer(even)
}输出:
输出: 4
输出: 16
输出: 363. Pipeline 的优点
- 关注点分离:每个阶段只做一件事,代码清晰。
- 可组合:阶段可以自由串联、重组。
- 流水线并行:不同阶段可以同时处理不同数据,提高吞吐量。
4. 注意:每个阶段都要关闭输出 channel
go
package main
import "fmt"
// ✅ 正确:每个阶段都 close 自己的输出
func stage(name string, in <-chan int) <-chan int {
out := make(chan int)
go func() {
defer close(out) // 关键:关闭输出,让下游 range 能结束
for v := range in {
out <- v + 1
}
fmt.Println(name, "阶段结束")
}()
return out
}
func main() {
ch := make(chan int)
go func() {
defer close(ch)
for i := 1; i <= 3; i++ {
ch <- i
}
}()
out := stage("A", ch)
out = stage("B", out)
out = stage("C", out)
for v := range out {
fmt.Println(v)
}
}如果不关闭,下游的 for range 会永久阻塞,goroutine 泄漏。
三、Fan-out / Fan-in 模式
Fan-out/Fan-in 是 Pipeline 的增强:Fan-out 把一个 stage 的输出分发给多个 goroutine 并行处理,Fan-in 把多个 goroutine 的输出合并到一个 channel。
1. 模式结构
┌→ worker1 ─┐
输入 → 分发 ─→ worker2 ─→ 合并 → 结果
└→ worker3 ─┘2. 完整实现
go
package main
import (
"fmt"
"sync"
)
// 生成器
func generate(nums ...int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for _, n := range nums {
out <- n
}
}()
return out
}
// 单个 worker:平方
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
}
// Fan-out:启动多个 square,共享同一个输入
// Fan-in:合并多个输出
func merge(cs ...<-chan int) <-chan int {
var wg sync.WaitGroup
out := make(chan int)
// 为每个输入 channel 启动一个转发 goroutine
for _, c := range cs {
wg.Add(1)
go func(ch <-chan int) {
defer wg.Done()
for v := range ch {
out <- v
}
}(c)
}
// 所有转发完成后关闭输出
go func() {
wg.Wait()
close(out)
}()
return out
}
func main() {
in := generate(1, 2, 3, 4, 5, 6, 7, 8, 9, 10)
// Fan-out:启动 3 个 square 并行处理
c1 := square(in)
c2 := square(in)
c3 := square(in)
// Fan-in:合并 3 个输出
for result := range merge(c1, c2, c3) {
fmt.Println(result)
}
}注意:这里 3 个 square 共享同一个 in channel,每个数据只会被其中一个 worker 取走(channel 是排他的)。这实现了负载分担。
3. Fan-out 何时有用
- 单个 stage 是性能瓶颈:用多个 worker 并行加速。
- 任务之间无依赖:可以独立处理。
- I/O 密集:多个 worker 同时等待 I/O,提高利用率。
四、Generator 模式
Generator(生成器)模式:用 goroutine + channel 生成数据流,实现「惰性求值」。消费者取一个,生成器才产一个。
1. 基本生成器
go
package main
import "fmt"
// 生成自然数序列
func naturals() <-chan int {
out := make(chan int)
go func() {
defer close(out)
for i := 1; ; i++ {
out <- i
}
}()
return out
}
func main() {
ch := naturals()
// 只取前 5 个
for i := 0; i < 5; i++ {
fmt.Println(<-ch)
}
// 注意:生成器 goroutine 会泄漏(永远在等接收)
// 实际应用应配合 context 或 done channel
}2. 带取消的生成器
go
package main
import (
"context"
"fmt"
"time"
)
// 带取消的斐波那契生成器
func fibonacci(ctx context.Context) <-chan int {
out := make(chan int)
go func() {
defer close(out)
a, b := 0, 1
for {
select {
case <-ctx.Done():
return
case out <- a:
a, b = b, a+b
}
}
}()
return out
}
func main() {
ctx, cancel := context.WithCancel(context.Background())
defer cancel() // 确保生成器能退出
ch := fibonacci(ctx)
for v := range ch {
fmt.Println(v)
if v > 100 {
cancel() // 取消生成器
break
}
}
// 给生成器一点时间退出
time.Sleep(100 * time.Millisecond)
fmt.Println("结束")
}用 select 监听 ctx.Done() 是避免 goroutine 泄漏的关键。
3. 生成器的应用
- 无限序列:自然数、斐波那契、素数。
- 流式数据:读取文件行、消息队列。
- 惰性求值:只生成需要的数据。
五、Future / Promise 模式
Future/Promise 模式:异步执行一个耗时操作,立即返回一个「未来的结果」句柄,需要结果时再等待。
1. 基本实现
go
package main
import (
"fmt"
"time"
)
// Future 表示一个未来的结果
type Future struct {
result chan int
}
// 异步执行函数,返回 Future
func asyncCompute(n int) *Future {
f := &Future{result: make(chan int, 1)}
go func() {
time.Sleep(time.Second) // 模拟耗时计算
f.result <- n * n
}()
return f
}
// Get 阻塞等待结果
func (f *Future) Get() int {
return <-f.result
}
func main() {
start := time.Now()
// 发起异步计算,立即返回
future := asyncCompute(42)
// 这期间可以做别的事
fmt.Println("做其他事情...")
time.Sleep(500 * time.Millisecond)
// 需要结果时阻塞等待
value := future.Get()
fmt.Println("结果:", value)
fmt.Println("总耗时:", time.Since(start)) // 约 1 秒(不是 1.5 秒)
}2. 带超时的 Future
go
package main
import (
"fmt"
"time"
)
type Result struct {
Value int
Err error
}
func asyncCompute(n int, timeout time.Duration) <-chan Result {
out := make(chan Result, 1)
go func() {
// 模拟可能超时的计算
time.Sleep(time.Second)
out <- Result{Value: n * 2}
}()
return out
}
func main() {
result := asyncCompute(42, 500*time.Millisecond)
select {
case r := <-result:
fmt.Println("结果:", r.Value)
case <-time.After(500 * time.Millisecond):
fmt.Println("超时")
}
}3. 并发 Future + 等待全部
go
package main
import (
"fmt"
"sync"
"time"
)
func asyncTask(name string, delay time.Duration, wg *sync.WaitGroup, result chan<- string) {
defer wg.Done()
time.Sleep(delay)
result <- fmt.Sprintf("%s 用时 %v", name, delay)
}
func main() {
start := time.Now()
var wg sync.WaitGroup
result := make(chan string, 3)
// 并发执行 3 个任务
wg.Add(3)
go asyncTask("任务A", 1*time.Second, &wg, result)
go asyncTask("任务B", 2*time.Second, &wg, result)
go asyncTask("任务C", 3*time.Second, &wg, result)
// 等待全部完成
go func() {
wg.Wait()
close(result)
}()
for r := range result {
fmt.Println(r)
}
fmt.Println("总耗时:", time.Since(start)) // 约 3 秒(取最长的)
}六、Timeout 模式
超时模式确保 goroutine 不会无限等待,是健壮并发程序必备的。我们在 select 章节已经接触过,这里做系统总结。
1. 单次操作超时
go
package main
import (
"fmt"
"time"
)
func slowDBQuery() <-chan string {
ch := make(chan string, 1)
go func() {
time.Sleep(2 * time.Second)
ch <- "数据"
}()
return ch
}
func main() {
select {
case res := <-slowDBQuery():
fmt.Println("查询成功:", res)
case <-time.After(500 * time.Millisecond):
fmt.Println("查询超时")
}
}2. 循环操作超时
go
package main
import (
"fmt"
"time"
)
func main() {
timeout := time.After(2 * time.Second)
ticker := time.NewTicker(300 * time.Millisecond)
defer ticker.Stop()
count := 0
loop:
for {
select {
case <-ticker.C:
count++
fmt.Printf("处理 %d\n", count)
case <-timeout:
fmt.Println("总超时,退出")
break loop
}
}
}3. 每个操作独立超时
go
package main
import (
"fmt"
"math/rand"
"time"
)
func process(id int) <-chan string {
ch := make(chan string, 1)
go func() {
delay := time.Duration(rand.Intn(3000)) * time.Millisecond
time.Sleep(delay)
ch <- fmt.Sprintf("任务 %d 完成 (耗时 %v)", id, delay)
}()
return ch
}
func main() {
rand.Seed(time.Now().UnixNano())
for i := 1; i <= 5; i++ {
select {
case res := <-process(i):
fmt.Println(res)
case <-time.After(time.Second):
fmt.Printf("任务 %d 超时\n", i)
}
}
}七、Cancel 模式
Cancel 模式:通过一个 channel 广播取消信号,所有监听它的 goroutine 都能收到并退出。
1. 基本 Cancel
go
package main
import (
"fmt"
"sync"
"time"
)
func worker(id int, cancel <-chan struct{}, wg *sync.WaitGroup) {
defer wg.Done()
for {
select {
case <-cancel:
fmt.Printf("worker %d 被取消\n", id)
return
default:
fmt.Printf("worker %d 工作中\n", id)
time.Sleep(300 * time.Millisecond)
}
}
}
func main() {
cancel := make(chan struct{})
var wg sync.WaitGroup
// 启动 3 个 worker
for i := 1; i <= 3; i++ {
wg.Add(1)
go worker(i, cancel, &wg)
}
// 1.5 秒后取消所有
time.Sleep(1500 * time.Millisecond)
close(cancel) // 广播取消
wg.Wait()
fmt.Println("全部退出")
}close(cancel) 会广播——所有从 cancel 接收的 goroutine 都会立即收到零值,从而退出。
2. Cancel + 数据 channel
go
package main
import (
"fmt"
"sync"
"time"
)
func processor(id int, jobs <-chan int, cancel <-chan struct{}, wg *sync.WaitGroup) {
defer wg.Done()
for {
select {
case <-cancel:
fmt.Printf("processor %d 取消\n", id)
return
case job, ok := <-jobs:
if !ok {
return
}
fmt.Printf("processor %d 处理 job %d\n", id, job)
time.Sleep(200 * time.Millisecond)
}
}
}
func main() {
jobs := make(chan int, 10)
cancel := make(chan struct{})
var wg sync.WaitGroup
for i := 1; i <= 2; i++ {
wg.Add(1)
go processor(i, jobs, cancel, &wg)
}
// 投递任务
for i := 1; i <= 5; i++ {
jobs <- i
}
// 提前取消
time.Sleep(300 * time.Millisecond)
close(cancel)
wg.Wait()
fmt.Println("结束")
// jobs 中可能还有未处理的任务,但 worker 已退出
}八、模式对比与选择建议
| 模式 | 适用场景 | 关键工具 |
|---|---|---|
| Worker Pool | 控制并发数、批量任务 | 固定 goroutine + chan |
| Pipeline | 多阶段数据处理 | 串联的 stage + chan |
| Fan-out/Fan-in | 并行加速单阶段、合并多源 | 多 goroutine + merge |
| Generator | 生成数据流、惰性求值 | goroutine + chan |
| Future/Promise | 异步操作、立即返回句柄 | chan + 句柄类型 |
| Timeout | 防止无限等待 | select + time.After |
| Cancel | 广播退出信号 | close(done chan) |
选择思路
- 需要控制并发数? → Worker Pool
- 任务有多个处理阶段? → Pipeline
- 某个阶段太慢需要并行? → Fan-out/Fan-in
- 需要生成数据流? → Generator
- 异步操作拿未来结果? → Future/Promise
- 怕操作卡死? → Timeout
- 需要通知所有 goroutine 退出? → Cancel
这些模式经常组合使用:比如 Pipeline 的某个 stage 用 Fan-out 加速,整个 pipeline 用 Cancel 控制退出,每个操作用 Timeout 兜底。
九、综合实战:可取消的并行数据处理流水线
把多种模式组合起来,实现一个真实场景:并发处理一批数据,支持取消和超时。
go
package main
import (
"context"
"fmt"
"math/rand"
"sync"
"time"
)
// 数据项
type Item struct {
ID int
Value int
}
// 处理结果
type Outcome struct {
ItemID int
Result int
Err error
}
// 阶段1:生成数据
func generate(ctx context.Context, count int) <-chan Item {
out := make(chan Item)
go func() {
defer close(out)
for i := 1; i <= count; i++ {
select {
case <-ctx.Done():
return
case out <- Item{ID: i, Value: rand.Intn(100)}:
}
}
}()
return out
}
// 阶段2:单个处理函数
func processOne(ctx context.Context, item Item) Outcome {
delay := time.Duration(rand.Intn(300)) * time.Millisecond
select {
case <-time.After(delay):
if item.Value < 10 {
return Outcome{ItemID: item.ID, Err: fmt.Errorf("值太小: %d", item.Value)}
}
return Outcome{ItemID: item.ID, Result: item.Value * item.Value}
case <-ctx.Done():
return Outcome{ItemID: item.ID, Err: ctx.Err()}
}
}
// 阶段2 + Fan-out:N 个 worker 并行处理
func process(ctx context.Context, in <-chan Item, numWorkers int) <-chan Outcome {
out := make(chan Outcome)
var wg sync.WaitGroup
for w := 1; w <= numWorkers; w++ {
wg.Add(1)
go func(id int) {
defer wg.Done()
for item := range in {
select {
case <-ctx.Done():
return
case out <- processOne(ctx, item):
}
}
}(w)
}
go func() {
wg.Wait()
close(out)
}()
return out
}
// 阶段3:收集结果
func collect(ctx context.Context, in <-chan Outcome) {
var success, failed int
for {
select {
case <-ctx.Done():
fmt.Printf("收集被取消: 成功 %d, 失败 %d\n", success, failed)
return
case o, ok := <-in:
if !ok {
fmt.Printf("完成: 成功 %d, 失败 %d\n", success, failed)
return
}
if o.Err != nil {
failed++
fmt.Printf("[失败] item %d: %v\n", o.ItemID, o.Err)
} else {
success++
fmt.Printf("[成功] item %d → %d\n", o.ItemID, o.Result)
}
}
}
}
func main() {
rand.Seed(time.Now().UnixNano())
// 5 秒超时
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
// 构建流水线:generate(20) → process(4 workers) → collect
items := generate(ctx, 20)
outcomes := process(ctx, items, 4)
collect(ctx, outcomes)
}这个例子综合了:Generator(generate)、Fan-out(4 个 worker)、Cancel(context)、Timeout(每项处理超时 + 总超时)、Pipeline(三阶段串联)。这是生产环境并发处理的典型骨架。
十、小结
本篇我们学习了 Go 并发编程的经典模式:
Worker Pool:固定数量 worker 处理任务,控制并发数。结构:任务 chan → N worker → 结果 chan。优点是并发可控、goroutine 复用。
Pipeline:多阶段串联,每阶段一个 goroutine,用 chan 连接。优点是关注点分离、流水线并行。注意每阶段要 close 输出。
Fan-out/Fan-in:Fan-out 分发给多个 worker 并行,Fan-in 合并结果。用于并行加速瓶颈阶段。
Generator:goroutine + chan 生成数据流,实现惰性求值。务必配合 cancel 机制避免泄漏。
Future/Promise:异步操作返回未来结果句柄,需要时 Get。支持并发执行 + 等待。
Timeout:select + time.After 防止无限等待。可单次超时、循环超时、每项独立超时。
Cancel:close(done chan) 广播退出信号,所有监听者都能收到。
模式选择:控制并发用 Worker Pool,多阶段用 Pipeline,并行加速用 Fan-out/Fan-in,生成数据用 Generator,异步结果用 Future,防卡死用 Timeout,广播退出用 Cancel。这些模式常组合使用。
下一篇我们将学习 context 包,它是 Go 官方提供的「取消、超时、传值」统一方案,能让上面这些模式的取消和超时更标准化。