goroutine + channel 是 Go 的招牌,但会用和用好是两码事。这篇整理生产环境最实用的并发模式——每个都带完整可运行的代码。


Q1:Pipeline 模式是什么?

A: 数据经过多个阶段处理,每阶段一个 goroutine,用 channel 连接:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
// 阶段1:生成数字
func generate(nums ...int) <-chan int {
out := make(chan int)
go func() {
for _, n := range nums {
out <- n
}
close(out)
}()
return out
}

// 阶段2:平方
func square(in <-chan int) <-chan int {
out := make(chan int)
go func() {
for n := range in {
out <- n * n
}
close(out)
}()
return out
}

// 阶段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
}

// 使用 —— 像搭积木一样组合
func main() {
// 1 → 4 → 9 → 16 → 25 → 过滤 >10 → [16, 25]
result := filter(square(generate(1, 2, 3, 4, 5)), 10)
for n := range result {
fmt.Println(n) // 16, 25
}
}

关键点:每个阶段的输出 channel 要 close(),否则下游的 range 会死锁。

Q2:Fan-Out / Fan-In 是什么?什么时候用?

A:

Fan-Out:一个输入分给多个 worker

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
func fanOut(in <-chan int, workers int) []<-chan int {
outs := make([]<-chan int, workers)
for i := 0; i < workers; i++ {
outs[i] = worker(in) // 每个 worker 独立消费 input
}
return outs
}

func worker(in <-chan int) <-chan int {
out := make(chan int)
go func() {
for n := range in {
// 模拟耗时操作
result := heavyComputation(n)
out <- result
}
close(out)
}()
return out
}

Fan-In:多个输出合并成一个

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
func fanIn(chans ...<-chan int) <-chan int {
var wg sync.WaitGroup
merged := make(chan int)

for _, ch := range chans {
wg.Add(1)
go func(c <-chan int) { // 注意参数传入!
defer wg.Done()
for v := range c {
merged <- v
}
}(ch)
}

go func() {
wg.Wait()
close(merged) // 所有输入都关闭后才关闭合并后的 channel
}()

return merged
}

典型场景:并行请求多个微服务,然后合并结果:

1
2
3
4
5
6
7
8
9
// 同时查询用户信息、订单、积分
userCh := fetchUser(ctx, userID)
orderCh := fetchOrders(ctx, userID)
pointsCh := fetchPoints(ctx, userID)

// 合并结果
for result := range fanIn(userCh, orderCh, pointsCh) {
process(result)
}

Q3:Worker Pool(工作池)—— 最实用的模式

A: 这是你 game-bass 和 slg-go 项目里最可能用到的模式:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
type Task struct {
ID int
Payload interface{}
}

func ProcessWithWorkerPool(tasks []Task, numWorkers int) []result {
var wg sync.WaitGroup

tasksCh := make(chan Task, len(tasks))
resultsCh := make(chan result, len(tasks))

// 启动固定数量的 worker
for w := 0; w < numWorkers; w++ {
wg.Add(1)
go func(workerID int) {
defer wg.Done()
for task := range tasksCh {
fmt.Printf("[Worker %d] 处理任务 #%d\n", workerID, task.ID)
r := doWork(task)
resultsCh <- r
}
}(w)
}

// 分发任务
for _, task := range tasks {
tasksCh <- task
}
close(tasksCh) // 重要!通知 worker 没有更多任务了

// 等待所有 worker 完成,然后关掉 results
go func() {
wg.Wait()
close(resultsCh)
}()

// 收集结果
var results []result
for r := range resultsCh {
results = append(results, r)
}

return results
}

为什么需要 Worker Pool?

无限制 goroutine Worker Pool
10 万个请求 → 10 万个 goroutine 固定 N 个 worker
内存爆炸、CPU 抢占 资源可控、吞吐稳定
OOM 风险 安全

实际项目中的变体——带限速的 Worker Pool:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
type RateLimitedPool struct {
workers int
rateLimit <-chan time.Time // 用 ticker 做限速
taskCh chan Task
resultCh chan result
}

func NewRateLimitedPool(workers int, qps int) *RateLimitedPool {
return &RateLimitedPool{
workers: workers,
rateLimit: time.Tick(time.Second / time.Duration(qps)),
taskCh: make(chan Task),
resultCh: make(chan result),
}
}

Q4:Semaphore(信号量)怎么实现?控制最大并发数

A: 用带缓冲的 channel 就能做信号量:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
func LimitedDo(tasks []Task, maxConcurrent int) {
sem := make(chan struct{}, maxConcurrent) // 信号量

var wg sync.WaitGroup
for _, task := range tasks {
wg.Add(1)
go func(t Task) {
defer wg.Done()
sem <- struct{}{} // 获取令牌(满了就阻塞)
defer func() { <-sem }() // 释放令牌

doWork(t)
}(task)
}
wg.Wait()
}

应用场景:

  • 限制数据库连接数
  • 控制 API 调用频率
  • 批量 HTTP 请求但不超过 N 个同时进行

Q5:Select 怎么用才算到位?

A: select 不只是”等待多个 channel”,它是 Go 并发的瑞士军刀:

基础用法:多路复用

1
2
3
4
5
6
7
8
select {
case data := <-dataCh:
handleData(data)
case err := <-errCh:
logError(err)
case <-time.After(3 * time.Second):
log.Println("超时了")
}

进阶用法:非阻塞收发

1
2
3
4
5
6
select {
case resultCh <- result:
fmt.Println("发送成功")
default:
fmt.Println("channel 满,稍后重试")
}

进阶用法:循环退出

1
2
3
4
5
6
7
8
9
10
11
12
13
14
func worker(ctx context.Context, jobs <-chan Job) {
for {
select {
case job, ok := <-jobs:
if !ok {
return // channel 已关闭,退出
}
process(job)
case <-ctx.Done():
fmt.Println("收到取消信号:", ctx.Err())
return // 优雅退出
}
}
}

⚠️ 注意:select 的随机性

如果多个 case 同时就绪,Go 随机选择一个执行。这不是 bug,是设计——防止某个 channel 总是被优先处理导致饥饿。

Q6:goroutine 泄漏怎么检测和避免?

A: 泄漏 = goroutine 永远不会退出。常见原因和解决方案:

泄漏原因 示例 解决方案
忘记 close channel sender 关了但 receiver 在等 用 context 取消或确保 close
select 缺少 default/timeout 阻塞在永远没数据的 channel 加超时或 done channel
goroutine 里用了阻塞调用且无法取消 http.Get 无超时 用 context.WithTimeout
WaitGroup 计数错误 Add/Done 不配对 确保 Add(1) 有对应 Done()

检测工具:

1
2
3
4
5
6
7
# 运行时查看 goroutine 数量
# 在代码中:
fmt.Printf("goroutine count: %d\n", runtime.NumGoroutine())

# pprof 分析
import _ "net/http/pprof"
# 访问: http://localhost:6060/debug/pprof/goroutine?debug=2

下期预告

下一篇深入 sync 包——Mutex、RWMutex、WaitGroup、Once、Pool。选对锁类型能显著提升性能,选错了就是隐蔽的性能杀手。然后最后一篇我们把所有知识串起来:从零搭建一个完整的 HTTP 服务。

TODO: 看看你 slg-go 项目里的 goroutine 使用——有没有哪个地方在 for 循环里启动 goroutine 但没有 context 控制取消?有没有应该用 Worker Pool 却直接裸起 goroutine 的地方?