当前位置: 首页 > news >正文

并发模式: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 模式,把多个处理阶段串成流水线,实现流式数据处理。

http://www.jsqmd.com/news/1395266/

相关文章:

  • 一条命令,把 M3U8 在线视频完整下载到本地:m3u8-downloader 快速上手
  • 【SSH 深度连载第 1 篇】SSH 是什么?为什么它是运维人员的 “空气和水“?从协议本质到替代方案全景解析
  • 2026国产压缩软件排行榜:功能与实用性综合评估排序
  • 论文查重率高的原因与智能降重实战技巧
  • 2026年8月郑州市新郑市联通500M宽带实测办理全流程 - 找卡家园
  • 网站离线下载保姆级攻略:WebSite-Downloader 整站保存实战
  • 基于MiniMax H3 LoRA训练器实现大模型高效微调与个性化定制
  • Maven依赖爆红全解析:从环境配置到依赖冲突的完整解决指南
  • Codex免费一键接入器:部署、测试与常见问题全解析
  • NVIDIA Profile Inspector 完整上手手册:免费解锁驱动隐藏设置的 4 个关键操作
  • 基于PatchCore的AMOLED Mura缺陷无监督检测实战指南
  • 深入解析Linux系统启动流程:从UEFI固件到systemd服务管理
  • 自托管服务实战(3):自建密码管理与文件同步
  • 数学建模竞赛B题全攻略:从破题到论文撰写的深度解析
  • 一套键鼠控制所有电脑:Input Leap 开源软件 KVM 完整实战指南
  • 2026国内主流压缩软件排行:按使用场景精准排序
  • 180、飞控中的多机协同:任务分配与调度
  • Gin进阶:参数绑定、验证与文件上传
  • CSDN 付费专栏|为什么要算 5G 真空波长与波数?从一道简单计算题看懂射频工程底层逻辑
  • 基于LP3667B的5V/1A反激式开关电源设计全解析
  • 老iPhone焕发新生机:palera1n越狱工具零基础实操攻略
  • 手机屏幕秒上电脑:scrcpy完整投屏与反向操控指南
  • YOLO室内手势识别手目标检测数据集-10896张
  • 心电图学习指南:从贺银成视频到高效笔记构建与临床诊断思维
  • Ace Data Cloud 定时任务:让 AI 自动搜热点、写文章、配图并发布
  • 不搞明白这10件事,还敢外包PCB设计?
  • DeepSeek Harness 本地搭建部署教程:接入 OpenAI 兼容 API
  • 锐捷交换机运维必备:十大核心查看命令与分层排查实战
  • 数据库实战:从SQL语法到索引优化与性能调优的体系化训练
  • 低功耗高性能笔记本推荐:2026年移动办公与创作的理性之选