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

相关新闻

网页爬虫的法律边界与合规技术实践

网页爬虫的法律边界与合规技术实践

1. 网页爬取的法律边界与行业现状 第一次接触网页爬取技术时,我和大多数开发者一样,最关心的不是技术实现,而是"这玩意儿到底能不能用"。2019年某电商平台起诉爬虫开发者的案件,让整个技术圈都开始重新审视这个灰色地带…

2026/8/15 2:05:05 阅读更多 →
Oracle分页查询性能优化:从ROWNUM原理到千万级数据实战

Oracle分页查询性能优化:从ROWNUM原理到千万级数据实战

1. 从一次深夜告警说起:为什么分页查询不是小事那天晚上十一点,我正打算关电脑,突然收到监控系统的告警,提示某个核心业务接口的响应时间飙升到了5秒以上。登录服务器一看,CPU和内存都还正常,但数据库的活跃…

2026/8/15 2:05:05 阅读更多 →
如何用 BoneAnimCopy 在 Blender 中一步到位复制骨骼动画:新手完整指南

如何用 BoneAnimCopy 在 Blender 中一步到位复制骨骼动画:新手完整指南

如何用 BoneAnimCopy 在 Blender 中一步到位复制骨骼动画:新手完整指南 【免费下载链接】blender_BoneAnimCopy 用于在blender中桥接骨骼动画的插件 项目地址: https://gitcode.com/gh_mirrors/bl/blender_BoneAnimCopy 做独立游戏或短片动画时,你…

2026/8/15 2:05:05 阅读更多 →

最新新闻

Typora图片处理全攻略:从基础插入到HTML/CSS高级美化

Typora图片处理全攻略:从基础插入到HTML/CSS高级美化

1. 从“拖拽”到“精修”:为什么图片处理是Markdown写作的最后一公里如果你用过Typora,大概率会和我一样,被它那种“所见即所得”的Markdown编辑体验所吸引。写标题、列清单、加粗文字,一切都流畅得像在Word里操作,但又…

2026/8/15 2:51:17 阅读更多 →
解决Chrome/Edge扩展无法启用:从.crx文件失效到解包安装全攻略

解决Chrome/Edge扩展无法启用:从.crx文件失效到解包安装全攻略

1. 问题场景:当你的浏览器扩展变成“幽灵”作为一名常年与各种浏览器插件打交道的开发者,我敢说,几乎每个深度用户都遇到过这个让人瞬间血压升高的场景:你从某个可靠的开发者论坛、技术社区或者朋友那里,拿到了一个功能…

2026/8/15 2:51:17 阅读更多 →
Mac开发者必备:Homebrew包管理器核心原理与高效使用指南

Mac开发者必备:Homebrew包管理器核心原理与高效使用指南

1. 项目概述:为什么Mac用户离不开Homebrew 如果你是一名Mac开发者,或者只是想在Mac上装点软件,那你大概率听说过Homebrew,也就是大家常说的 brew 。它远不止是一个简单的安装命令,而是整个macOS生态里,一…

2026/8/15 2:51:17 阅读更多 →
IDEA翻译插件配置指南:集成百度翻译API提升开发效率

IDEA翻译插件配置指南:集成百度翻译API提升开发效率

1. 项目缘起:为什么我们需要一个靠谱的翻译插件?作为一名开发者,我每天都要和大量的英文文档、技术博客、开源代码注释以及报错信息打交道。相信很多同行都有过类似的经历:在IntelliJ IDEA里写代码,突然遇到一个陌生的…

2026/8/15 2:51:17 阅读更多 →
MySQL从安装到精通:避坑指南与核心概念全解析

MySQL从安装到精通:避坑指南与核心概念全解析

1. 从“安装即放弃”到“丝滑上手”:一个老DBA的MySQL避坑心法我见过太多新手,兴致勃勃地下载了MySQL,结果在安装配置的第一步就卡壳,要么服务起不来,要么连不上,最后只能无奈放弃,转头去找那些…

2026/8/15 2:51:17 阅读更多 →
[通信与计算]微积分:基础概念及其在通信中的应用

[通信与计算]微积分:基础概念及其在通信中的应用

本文从工程实践角度介绍微积分基础,包括极限、导数、积分和Taylor级数,并说明这些概念在通信和信号处理系统分析与设计中的具体用法。图1:函数f(x)sin(x)0.5x及其在x00.5处的切线示意,将导数解释为斜率。图2:f(x)e^{-0…

2026/8/15 2:50:17 阅读更多 →

日新闻

内景 空间站内部 中国空间站 太空 内仓

内景 空间站内部 中国空间站 太空 内仓

本项目为前几天收费帮学妹做的一个项目,在工作环境中基本使用不到,但是很多学校把这个当作编程入门的项目来做,故分享出本项目供初学者参考。 一、项目描述 空间站内部 中国空间站 太空 内仓 地址:本地PC端运行(或Web…

2026/8/15 0:00:30 阅读更多 →
重新定义数据接口:3个突破性场景让通达信数据读取更智能

重新定义数据接口:3个突破性场景让通达信数据读取更智能

重新定义数据接口:3个突破性场景让通达信数据读取更智能 【免费下载链接】mootdx 通达信数据读取的一个简便使用封装 项目地址: https://gitcode.com/GitHub_Trending/mo/mootdx 当我们面对海量金融数据时,传统的数据获取方式往往让我们陷入困境—…

2026/8/15 0:00:30 阅读更多 →
一文读懂快消WMS怎么选?2026年国内外10大主流WMS品牌盘点

一文读懂快消WMS怎么选?2026年国内外10大主流WMS品牌盘点

快消品(FMCG)是流通速度较快、竞争较为激烈的行业之一。一瓶饮料从出厂到消费者手中,往往只有几十天甚至几天的周转窗口。这决定了快消行业的仓储管理系统(WMS)与制造业、电商行业存在明显区别:它不仅需要管…

2026/8/15 0:02:30 阅读更多 →

周新闻

5分钟告别提取码焦虑:baidupankey如何智能破解百度网盘资源锁

5分钟告别提取码焦虑:baidupankey如何智能破解百度网盘资源锁

5分钟告别提取码焦虑:baidupankey如何智能破解百度网盘资源锁 【免费下载链接】baidupankey 在线查询网盘提取码(维护中 rm repo) 项目地址: https://gitcode.com/gh_mirrors/ba/baidupankey 你是否曾经在深夜寻找一份重要资料&#x…

2026/8/13 2:38:34 阅读更多 →
如何快速生成中国车牌图片:Python开源工具完整指南

如何快速生成中国车牌图片:Python开源工具完整指南

如何快速生成中国车牌图片:Python开源工具完整指南 【免费下载链接】chinese_license_plate_generator 中国车牌生成器 项目地址: https://gitcode.com/gh_mirrors/ch/chinese_license_plate_generator 中国车牌生成器是一个基于Python的开源项目&#xff0c…

2026/8/13 10:41:52 阅读更多 →
收藏!小白程序员轻松入门大模型,从Harness工程开始实践

收藏!小白程序员轻松入门大模型,从Harness工程开始实践

文章强调学习大模型不应只关注模型本身,而应重视模型外的系统搭建,即Harness。提出AgentModelHarness的实用公式,详细介绍Harness的四个层次:持久化层、执行层、控制层和观察与验证层。文章还探讨了上下文工程、工具设计、AGENTS.…

2026/8/13 10:41:51 阅读更多 →

月新闻

免费解锁百度网盘SVIP加速:macOS用户必备的下载提速终极指南

免费解锁百度网盘SVIP加速:macOS用户必备的下载提速终极指南

免费解锁百度网盘SVIP加速:macOS用户必备的下载提速终极指南 【免费下载链接】BaiduNetdiskPlugin-macOS For macOS.百度网盘 破解SVIP、下载速度限制~ 项目地址: https://gitcode.com/gh_mirrors/ba/BaiduNetdiskPlugin-macOS 还在为百度网盘macOS版的龟速下…

2026/8/14 13:40:53 阅读更多 →
终极ncmdump指南:3分钟实现网易云NCM音乐解密与格式转换

终极ncmdump指南:3分钟实现网易云NCM音乐解密与格式转换

终极ncmdump指南:3分钟实现网易云NCM音乐解密与格式转换 【免费下载链接】ncmdump 项目地址: https://gitcode.com/gh_mirrors/ncmd/ncmdump 还在为网易云音乐下载的NCM格式文件无法在其他播放器播放而烦恼吗?ncmdump解密工具帮你轻松解决这个困…

2026/8/14 14:06:45 阅读更多 →
HarmonyOS 应用开发《掌上英语》第81篇: 智能体卡片:为英语学习 App 打造桌面级学习助手

HarmonyOS 应用开发《掌上英语》第81篇: 智能体卡片:为英语学习 App 打造桌面级学习助手

AgentCard 智能体卡片:为英语学习 App 打造桌面级学习助手适用平台:HarmonyOS 7.0 (API 26 Beta)一、引言 HarmonyOS 7.0(API 26 Beta)新增了 AgentCard 智能体卡片能力,这是继 HMAF(鸿蒙智能体框架&#x…

2026/8/15 2:35:29 阅读更多 →