第28篇 并发模式,Fan-in/Fan-out流水线
摘要:Fan-out将任务分发给多个goroutine并行处理,Fan-in将多个goroutine的结果汇聚到一个通道。本文从一个多数据源聚合场景说起,讲清楚扇出扇入的实现和通道合并的细节。
一个多数据源聚合需求
做过搜索系统的同学应该熟悉这个场景。用户搜一个关键词,后端要同时查商品库、店铺库、文章库、问答库,四个数据源的结果合并排序后返回。串行查的话每个数据源平均200毫秒,四个加起来800毫秒,用户体验很差。
并行查的话四个数据源同时跑,总耗时取决于最慢的那个,大概250毫秒,快了三倍。但问题来了,四个数据源各自返回一个结果通道,怎么把四个通道的结果合并成一个通道给下游消费?这就用到 Fan-out 和 Fan-in。
Fan-out 是扇出,把一个任务拆成多份分发给多个 goroutine 并行处理。Fan-in 是扇入,把多个 goroutine 的输出通道合并成一个通道。两个组合起来就是典型的分散计算、汇聚结果模式。
Fan-out 分发任务
先看 Fan-out,把一批搜索请求分发给多个数据源 worker 并行查询。
packagemainimport("fmt""math/rand""time")// SearchResult 搜索结果typeSearchResultstruct{Sourcestring// 数据源名称Items[]string// 搜索到的条目}// search 模拟在单个数据源中搜索funcsearch(source,keywordstring)SearchResult{// 模拟不同数据源的查询耗时,100到300毫秒随机delay:=time.Duration(100+rand.Intn(200))*time.Millisecond time.Sleep(delay)returnSearchResult{Source:source,Items:[]string{keyword+"-result-1",keyword+"-result-2"},}}// searchSource 在指定数据源搜索,结果写入通道funcsearchSource(source,keywordstring)<-chanSearchResult{out:=make(chanSearchResult,1)// 带缓冲,避免阻塞gofunc(){deferclose(out)// 查完关闭通道out<-search(source,keyword)// 在该数据源执行搜索}()returnout}// fanOut 把搜索任务分发给多个数据源并行查询// 返回每个数据源的结果通道funcfanOut(keywordstring,sources[]string)[]<-chanSearchResult{out:=make([]<-chanSearchResult,len(sources))fori,src:=rangesources{out[i]=searchSource(src,keyword)// 每个数据源一个goroutine}returnout}funcmain(){sources:=[]string{"商品库","店铺库","文章库","问答库"}start:=time.Now()// Fan-out: 四个数据源同时搜索channels:=fanOut("手机",sources)// 先简单收一下结果,下面用fan-in优雅处理for_,ch:=rangechannels{res:=<-ch// 阻塞等待每个数据源返回fmt.Printf("[%s] 找到 %d 条结果\n",res.Source,len(res.Items))}fmt.Printf("总耗时 %v\n",time.Since(start))}这里有个问题,上面的写法是逐个等结果,哪个数据源慢就要卡到最后。理想情况是哪个数据源先返回就先处理,这就需要 Fan-in 把多个通道合并。
Fan-in 合并结果
Fan-in 的核心是把多个输入通道合并成一个输出通道。做法是给每个输入通道起一个 goroutine,把结果转发到输出通道,所有 goroutine 完成后关闭输出通道。
packagemainimport("context""fmt""math/rand""sync""time")// SearchResult 搜索结果typeSearchResultstruct{SourcestringItems[]string}// search 模拟在单个数据源中搜索funcsearch(source,keywordstring)SearchResult{delay:=time.Duration(100+rand.Intn(200))*time.Millisecond time.Sleep(delay)returnSearchResult{Source:source,Items:[]string{keyword+"-r1",keyword+"-r2"},}}// searchSource 在指定数据源搜索,结果写入通道funcsearchSource(source,keywordstring)<-chanSearchResult{out:=make(chanSearchResult,1)gofunc(){deferclose(out)out<-search(source,keyword)}()returnout}// fanOut 分发给多个数据源并行查询funcfanOut(keywordstring,sources[]string)[]<-chanSearchResult{out:=make([]<-chanSearchResult,len(sources))fori,src:=rangesources{out[i]=searchSource(src,keyword)}returnout}// fanIn 合并多个输入通道为一个输出通道funcfanIn(channels[]<-chanSearchResult)<-chanSearchResult{varwg sync.WaitGroup out:=make(chanSearchResult,len(channels))// 缓冲足够大// 为每个输入通道启动一个转发goroutinefor_,ch:=rangechannels{wg.Add(1)gofunc(c<-chanSearchResult){deferwg.Done()forres:=rangec{// 读取直到通道关闭out<-res// 转发到合并通道}}(ch)}// 单独的goroutine等所有转发完成,然后关闭输出通道gofunc(){wg.Wait()close(out)// 所有输入读完才关闭,避免panic}()returnout}funcmain(){sources:=[]string{"商品库","店铺库","文章库","问答库"}start:=time.Now()// Fan-out: 分发给四个数据源并行搜索channels:=fanOut("手机",sources)// Fan-in: 合并四个结果通道为一个merged:=fanIn(channels)// 从合并通道消费结果,谁先返回谁先被处理forres:=rangemerged{fmt.Printf("[%s] 找到 %d 条结果\n",res.Source,len(res.Items))}fmt.Printf("总耗时 %v\n",time.Since(start))}现在不管哪个数据源先返回,都能立即被消费,不用等最慢的那个。总耗时接近最慢数据源的查询时间。close(out) 必须在单独的 goroutine 里执行,这一点很关键,下一节详细说。
独家踩坑,fanIn里goroutine泄漏
这个坑我在生产环境踩过。上面的 fanIn 看起来没问题,但如果调用方提前 break 了 for range merged 循环,比如找到足够结果就不再消费,那些转发 goroutine 就会阻塞在 out <- res 上,永远退不出。
// 泄漏场景,只消费前2个结果就退出merged:=fanIn(channels)count:=0forres:=rangemerged{fmt.Println(res.Source)count++ifcount>=2{break// 剩下的goroutine阻塞在 out <- res,泄漏了}}两个数据源的结果被消费了,但另外两个 goroutine 往 out 写入时阻塞,因为没人读了。out 通道有缓冲,但如果缓冲满了就卡住。这些 goroutine 永远不会退出,内存泄漏。
修复方案是引入 Context,让转发 goroutine 能感知取消信号。
// fanInCtx 带context的fan-in,支持取消funcfanInCtx(ctx context.Context,channels[]<-chanSearchResult)<-chanSearchResult{varwg sync.WaitGroup out:=make(chanSearchResult,len(channels))for_,ch:=rangechannels{wg.Add(1)gofunc(c<-chanSearchResult){deferwg.Done()for{select{caseres,ok:=<-c:if!ok{return// 输入通道关闭,退出}// 转发时也监听取消信号select{caseout<-res:// 正常转发case<-ctx.Done():return// 被取消,退出}case<-ctx.Done():return// 被取消,退出}}}(ch)}// 等所有转发goroutine完成再关闭输出gofunc(){wg.Wait()close(out)}()returnout}现在调用方 break 之前 cancel 一下 context,所有 goroutine 都能及时退出。这个嵌套 select 看着复杂,但逻辑很清晰,外层监听输入和取消,内层监听输出和取消。
对比分析
| 维度 | Fan-in/Fan-out | 串行查询 | WaitGroup并行 |
|---|---|---|---|
| 并发度 | 多数据源并行 | 1 | 多 |
| 结果顺序 | 谁快谁先 | 固定顺序 | 需等待全部 |
| 总耗时 | 接近最慢者 | 全部之和 | 接近最慢者 |
| 流式处理 | 支持 | 不支持 | 不支持 |
| 可取消 | 配合ctx | 难 | 中等 |
Fan-in/Fan-out 相比 WaitGroup 的优势在于流式处理。WaitGroup 要等所有 goroutine 完成才能拿到结果,Fan-in 可以谁先完成谁先处理,对用户体验更友好。
总结预告
Fan-out 分发任务实现并行,Fan-in 合并结果实现汇聚。两者组合是处理多数据源聚合的标准姿势。核心注意点是通道关闭时机和 goroutine 泄漏防护,Context 是防泄漏的利器。
下一篇讲 Pipeline 模式,把多个处理阶段串成流水线,实现流式数据处理。