Skip to content

并发模式: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
输出: 36

3. 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)

选择思路

  1. 需要控制并发数? → Worker Pool
  2. 任务有多个处理阶段? → Pipeline
  3. 某个阶段太慢需要并行? → Fan-out/Fan-in
  4. 需要生成数据流? → Generator
  5. 异步操作拿未来结果? → Future/Promise
  6. 怕操作卡死? → Timeout
  7. 需要通知所有 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 并发编程的经典模式:

  1. Worker Pool:固定数量 worker 处理任务,控制并发数。结构:任务 chan → N worker → 结果 chan。优点是并发可控、goroutine 复用。

  2. Pipeline:多阶段串联,每阶段一个 goroutine,用 chan 连接。优点是关注点分离、流水线并行。注意每阶段要 close 输出。

  3. Fan-out/Fan-in:Fan-out 分发给多个 worker 并行,Fan-in 合并结果。用于并行加速瓶颈阶段。

  4. Generator:goroutine + chan 生成数据流,实现惰性求值。务必配合 cancel 机制避免泄漏。

  5. Future/Promise:异步操作返回未来结果句柄,需要时 Get。支持并发执行 + 等待。

  6. Timeout:select + time.After 防止无限等待。可单次超时、循环超时、每项独立超时。

  7. Cancel:close(done chan) 广播退出信号,所有监听者都能收到。

  8. 模式选择:控制并发用 Worker Pool,多阶段用 Pipeline,并行加速用 Fan-out/Fan-in,生成数据用 Generator,异步结果用 Future,防卡死用 Timeout,广播退出用 Cancel。这些模式常组合使用。

下一篇我们将学习 context 包,它是 Go 官方提供的「取消、超时、传值」统一方案,能让上面这些模式的取消和超时更标准化。