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 func generate (nums ...int ) <-chan int { out := make (chan int ) go func () { for _, n := range nums { out <- n } close (out) }() return out } func square (in <-chan int ) <-chan int { out := make (chan int ) go func () { for n := range in { out <- n * n } close (out) }() return out } 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 () { result := filter(square(generate(1 , 2 , 3 , 4 , 5 )), 10 ) for n := range result { fmt.Println(n) } }
关键点 :每个阶段的输出 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) } 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) }() 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)) 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) 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 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 } 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 fmt.Printf("goroutine count: %d\n" , runtime.NumGoroutine()) import _ "net/http/pprof"
下期预告 下一篇深入 sync 包 ——Mutex、RWMutex、WaitGroup、Once、Pool。选对锁类型能显著提升性能,选错了就是隐蔽的性能杀手。然后最后一篇我们把所有知识串起来:从零搭建一个完整的 HTTP 服务 。
TODO : 看看你 slg-go 项目里的 goroutine 使用——有没有哪个地方在 for 循环里启动 goroutine 但没有 context 控制取消?有没有应该用 Worker Pool 却直接裸起 goroutine 的地方?