news 2026/8/15 2:05:48

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

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
并发模式: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("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 模式,把多个处理阶段串成流水线,实现流式数据处理。

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/8/15 2:05:30

红绿在黑白照片中的灰度差异解析

哈哈&#xff0c;这个问题问得非常棒&#xff0c;它触及了数字图像处理中一个非常经典且核心的概念——灰度化&#xff08;Grayscale&#xff09;。你的观察完全正确&#xff0c;但背后的原理比“完全分不清”要精妙得多。直接回答你的问题&#xff1a;在最终的黑白照片里&…

作者头像 李华
网站建设 2026/8/15 2:04:10

地球沉降危机:人类妄为加速维度崩塌037

摘要&#xff1a;本文以“维性力网”为核心概念&#xff0c;阐述天地万物依托本源之力构筑全域拓扑结构&#xff0c;并持续向下沉沦的宇宙观。文章指出&#xff0c;人为开挖山体、改造河湖等行为会在星球拓扑上留下淤堵伤痕&#xff0c;加速地球整体沉降&#xff0c;提前触发纪…

作者头像 李华
网站建设 2026/8/15 2:02:35

从提示工程到驾驭工程:构建可靠AI Agent的系统工程实践

1. 项目概述&#xff1a;从“提示”到“驾驭”的工程思维升级最近在AI Agent的开发圈子里&#xff0c;一个词的热度正在悄然攀升&#xff0c;那就是“Harness Engineering”&#xff0c;有人把它翻译成“驭缰工程”或“驾驭工程”。如果你和我一样&#xff0c;在过去一年里深陷…

作者头像 李华
网站建设 2026/8/15 1:58:06

BilibiliDown 新手完全指南:轻松搞定 B站视频下载与批量离线保存

BilibiliDown 新手完全指南&#xff1a;轻松搞定 B站视频下载与批量离线保存 【免费下载链接】BilibiliDown (GUI-多平台支持) B站 哔哩哔哩 视频下载器。支持稍后再看、收藏夹、UP主视频批量下载|Bilibili Video Downloader &#x1f633; 项目地址: https://gitcode.com/gh…

作者头像 李华
网站建设 2026/8/15 1:55:30

90%的人不知道:Windows右键菜单管理,原来3步就能搞定

90%的人不知道&#xff1a;Windows右键菜单管理&#xff0c;原来3步就能搞定 【免费下载链接】ContextMenuManager &#x1f5b1;️ 纯粹的Windows右键菜单管理程序 项目地址: https://gitcode.com/gh_mirrors/co/ContextMenuManager 你每天要右键点击多少次文件&#x…

作者头像 李华
网站建设 2026/8/15 1:54:57

GCC符号可见性详解:-fvisibility=hidden构建健壮动态库

1. 项目概述&#xff1a;为什么我们需要关心符号可见性&#xff1f;如果你在Linux或类Unix系统上写过C/C的动态库&#xff08;.so文件&#xff09;&#xff0c;或者为Windows平台写过DLL&#xff0c;那你大概率遇到过一些让人头疼的链接问题。比如&#xff0c;你精心编写的库&a…

作者头像 李华