ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

并发模式:Fan-in/Fan-out流水线

并发模式:Fan-in/Fan-out流水线 第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(fmtmath/randtime)// SearchResult 搜索结果typeSearchResultstruct{Sourcestring// 数据源名称Items[]string// 搜索到的条目}// search 模拟在单个数据源中搜索funcsearch(source,keywordstring)SearchResult{// 模拟不同数据源的查询耗时100到300毫秒随机delay:time.Duration(100rand.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(contextfmtmath/randsynctime)// SearchResult 搜索结果typeSearchResultstruct{SourcestringItems[]string}// search 模拟在单个数据源中搜索funcsearch(source,keywordstring)SearchResult{delay:time.Duration(100rand.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)countifcount2{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 模式把多个处理阶段串成流水线实现流式数据处理。
返回列表