Deprecated: imwpcache\f884414bce24ee67f\f73723ec7b1919fa5::__construct(): Implicitly marking parameter $YECBGYFECGEAFWHA as nullable is deprecated, the explicit nullable type must be used instead in /www/wwwroot/www.chuangxiangniao.com/wp-content/plugins/imwpcache-dist/build/f884414bce24ee67ff73723ec7b1919fa5.php on line 2

Deprecated: imwpcache\f884414bce24ee67f\f73723ec7b1919fa5::__construct(): Implicitly marking parameter $BBWFDDBHHYHDXXAB as nullable is deprecated, the explicit nullable type must be used instead in /www/wwwroot/www.chuangxiangniao.com/wp-content/plugins/imwpcache-dist/build/f884414bce24ee67ff73723ec7b1919fa5.php on line 2
Go并发模式:详解Fan-Out(一生产者多消费者)_创想鸟

Go并发模式:详解Fan-Out(一生产者多消费者)

Go并发模式:详解Fan-Out(一生产者多消费者)

本文深入探讨go语言中的fan-out并发模式,演示如何通过通道实现一生产者向多消费者分发数据副本。文章详细介绍了`fanout`函数的实现,包括创建缓冲通道以控制消费者滞后、数据分发协程的运作,以及在输入通道耗尽后正确关闭所有输出通道的关键机制,确保资源有效管理与并发流程的顺畅。

什么是Fan-Out模式?

在Go语言的并发编程中,”Fan-Out”(扇出)是一种常见的模式,它描述了一个单一的数据源(生产者)将数据分发给多个接收者(消费者)的场景。每个消费者都会接收到数据源发送的相同数据副本。这与”Fan-In”(扇入)模式相对,Fan-In模式是将多个生产者的输出汇聚到一个单一的消费者。

Fan-Out模式在许多场景下都非常有用,例如:

事件广播:当一个事件发生时,需要通知多个监听者。任务分发:一个主任务产生的数据需要被多个子任务并行处理。数据复制:将相同的数据流发送到不同的处理管道。

Go语言通过其强大的goroutine和channel机制,可以优雅地实现Fan-Out模式。

Fan-Out模式的核心实现

实现Fan-Out模式的关键在于创建一个中间层,它从一个输入通道读取数据,然后将这些数据的副本写入到多个输出通道。

函数签名与参数

一个典型的Fan-Out函数可以定义如下:

func fanOut(ch <-chan int, size, lag int) []chan int {    // ... 实现细节}

ch size int: 表示需要创建的输出通道的数量,即有多少个消费者将接收数据。lag int: 这是一个关键参数,用于控制输出通道的缓冲大小。它决定了消费者能够落后于生产者多少数据而不会阻塞整个系统。

创建输出通道

首先,我们需要根据size参数创建相应数量的输出通道。这些通道将构成一个切片返回给调用者。

    cs := make([]chan int, size)    for i, _ := range cs {        // lag参数决定了通道的缓冲大小        cs[i] = make(chan int, lag)    }

这里,make(chan int, lag) 创建了一个带有指定缓冲大小的通道。缓冲通道允许在发送者和接收者之间存在一定的数据量而不会立即阻塞。如果lag为0,则创建的是无缓冲通道。

数据分发协程

Fan-Out模式的核心是一个独立的goroutine,它负责从输入通道读取数据,并将每个数据副本发送到所有输出通道。

    go func() {        for i := range ch { // 从输入通道读取数据            for _, c := range cs { // 将数据副本发送到所有输出通道                c <- i            }        }        // 当输入通道关闭且所有数据被读取完毕后,关闭所有输出通道        for _, c := range cs {            close(c)        }    }()

这个goroutine会一直运行,直到输入通道ch被关闭且所有数据都被range循环读取完毕。

通道的关闭

通道的正确关闭是并发编程中非常重要的一环。在Fan-Out模式中,当输入通道ch耗尽(即生产者不再发送数据并关闭了它)时,Fan-Out协程应该关闭所有它创建的输出通道。

for i := range ch 循环会在ch关闭时自动退出。紧接着,for _, c := range cs { close(c) } 会遍历并关闭所有输出通道。这向所有消费者发出了信号,表明不再有新的数据到来,它们可以安全地退出循环或清理资源。

无缓冲通道的Fan-Out

作为对比,我们也可以实现一个使用无缓冲通道的Fan-Out函数:

func fanOutUnbuffered(ch <-chan int, size int) []chan int {    cs := make([]chan int, size)    for i, _ := range cs {        cs[i] = make(chan int) // 无缓冲通道    }    go func() {        for i := range ch {            for _, c := range cs {                c <- i            }        }        for _, c := range cs {            close(c)        }    }()    return cs}

与缓冲通道版本的主要区别在于make(chan int)。使用无缓冲通道意味着任何一个消费者如果未能及时接收数据,都将阻塞Fan-Out协程,进而阻塞所有其他输出通道的数据发送,甚至可能回溯到生产者。因此,在大多数实际应用中,推荐使用缓冲通道来提高系统的并发性和容错性。

示例代码解析

下面是一个完整的示例,演示了如何将生产者、Fan-Out函数和多个消费者组合起来。

package mainimport (    "fmt"    "time")// producer 函数模拟一个数据生产者// 它会生成指定数量的整数,并每秒发送一个func producer(iters int) <-chan int {    c := make(chan int)    go func() {        for i := 0; i < iters; i++ {            c <- i            time.Sleep(1 * time.Second) // 模拟生产耗时        }        close(c) // 生产完毕后关闭通道    }()    return c}// consumer 函数模拟一个数据消费者// 它从输入通道读取数据并打印func consumer(cin <-chan int) {    for i := range cin {        fmt.Printf("Consumer received: %dn", i)    }    fmt.Println("Consumer finished.")}// fanOut 函数实现带缓冲的Fan-Out模式// ch: 输入通道// size: 输出通道的数量// lag: 输出通道的缓冲大小func fanOut(ch <-chan int, size, lag int) []chan int {    cs := make([]chan int, size)    for i := range cs {        cs[i] = make(chan int, lag) // 创建带缓冲的输出通道    }    go func() {        for i := range ch { // 从输入通道读取数据            for _, c := range cs { // 将数据副本发送到所有输出通道                c <- i            }        }        // 输入通道关闭后,关闭所有输出通道        for _, c := range cs {            close(c)        }    }()    return cs}// fanOutUnbuffered 函数实现无缓冲的Fan-Out模式func fanOutUnbuffered(ch <-chan int, size int) []chan int {    cs := make([]chan int, size)    for i := range cs {        cs[i] = make(chan int) // 创建无缓冲的输出通道    }    go func() {        for i := range ch {            for _, c := range cs {                c <- i            }        }        for _, c := range cs {            close(c)        }    }()    return cs}func main() {    // 1. 创建一个生产者,生产10个数据    c := producer(10)    // 2. 使用fanOutUnbuffered函数创建3个输出通道    // 尝试将 fanOutUnbuffered 替换为 fanOut(c, 3, 1) 或 fanOut(c, 3, 5)    // 观察缓冲对行为的影响    chans := fanOutUnbuffered(c, 3)     // 3. 启动3个消费者    // 前两个消费者作为goroutine运行    go consumer(chans[0])    go consumer(chans[1])    // 最后一个消费者在主goroutine中运行,阻塞主goroutine直到其完成    consumer(chans[2])     fmt.Println("Main goroutine finished.")}

在main函数中:

producer(10) 创建了一个生产者,它将生成0到9的整数。fanOutUnbuffered(c, 3) (或 fanOut(c, 3, lag)) 将生产者的输出通道c分发给3个新的输出通道。go consumer(chans[0]), go consumer(chans[1]) 启动了两个并发的消费者。consumer(chans[2]) 在主goroutine中运行第三个消费者。这意味着主goroutine会等待这个消费者完成所有数据的接收和处理,这有助于确保所有goroutine在程序退出前有足够的时间运行。

注意事项与最佳实践

缓冲与阻塞:

无缓冲通道:如果任何一个消费者处理速度过慢,或者暂时未准备好接收数据,Fan-Out协程在尝试向该通道发送数据时会阻塞。这将导致所有其他输出通道的数据发送也暂停,甚至可能反向阻塞生产者。在需要严格同步或确保所有消费者同时处理数据的场景下可能适用,但通常会导致性能瓶颈。缓冲通道:通过lag参数设置合适的缓冲大小,可以允许消费者在一定程度上滞后于生产者和Fan-Out协程,而不会立即造成阻塞。这提高了系统的并发性和弹性。选择合适的缓冲大小需要根据实际应用场景进行权衡,过大的缓冲可能导致内存占用增加,过小的缓冲则可能仍然引起阻塞。

通道的关闭:

生产者关闭输入通道:生产者必须在所有数据发送完毕后关闭其输出通道。这是Fan-Out协程判断输入结束的信号。Fan-Out协程关闭输出通道:Fan-Out协程必须在输入通道关闭并处理完所有数据后,关闭所有它创建的输出通道。这向消费者发出信号,表明不再有数据到来,消费者可以安全退出for range循环。避免重复关闭:通道只能关闭一次,重复关闭会导致运行时panic。

错误处理:本示例为了简洁未包含错误处理。在实际应用中,生产者、Fan-Out协程和消费者都可能遇到错误。需要设计适当的错误传递机制(例如,通过额外的错误通道或结构体)来处理这些情况。

资源管理:确保所有启动的goroutine都能正常退出,避免goroutine泄漏。正确关闭通道是实现这一目标的关键。

总结

Fan-Out模式是Go语言并发编程中一个强大而灵活的工具,它使得一个生产者能够高效地将数据分发给多个消费者。通过合理利用Go的通道机制,特别是缓冲通道,我们可以构建出健壮、高性能的并发系统。理解缓冲通道的作用、数据分发协程的逻辑以及通道的正确关闭时机,是成功实现和应用Fan-Out模式的关键。

以上就是Go并发模式:详解Fan-Out(一生产者多消费者)的详细内容,更多请关注创想鸟其它相关文章!

版权声明:本文内容由互联网用户自发贡献,该文观点仅代表作者本人。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。
如发现本站有涉嫌抄袭侵权/违法违规的内容, 请发送邮件至 chuangxiangniao@163.com 举报,一经查实,本站将立刻删除。
发布者:程序猿,转转请注明出处:https://www.chuangxiangniao.com/p/1413212.html

赞 (0)
打赏 微信扫一扫 微信扫一扫 支付宝扫一扫 支付宝扫一扫
Go语言中包级变量的初始化顺序与依赖分析
上一篇 2025年12月16日 07:01:49
在 Go 语言中正确定义函数参数类型
下一篇 2025年12月16日 07:02:10

相关推荐

  • 使用EventBus实现Android实时速度显示与后台保存教程

    本教程详细介绍了如何在Android应用中实现实时速度的显示与后台保存功能。通过利用前台服务(Foreground Service)获取位置数据,并结合EventBus库实现服务与UI界面(MainActivity)之间的实时数据通信,确保即使应用处于后台或屏幕关闭时,速度数据也能持续更新并显示在用…

    2026年9月21日
    000
  • 提高蝴蝶号无人直播留存率的6个实用技巧和策略

    提高蝴蝶号无人直播留存率的6个实用技巧和策略提高蝴蝶号无人直播留存率的6个实用技巧和策略提高蝴蝶号无人直播留存率的6个实用技巧和策略提高蝴蝶号无人直播留存率的6个实用技巧和策略

    提高蝴蝶号无人直播留存率的核心在于让用户觉得直播间“有东西”,具体措施包括:1.内容为王,垂直深耕某一领域并提供专业知识;2.互动是魂,利用弹幕、投票、抽奖引导用户参与;3.利益驱动,通过抽奖、红包提升用户积极性;4.氛围营造,打造独特风格和专属互动方式;5.数据分析,持续优化直播策略;6.活动预告…

    2026年9月21日 • 用户投稿
    100
  • 佳能EOS R1对决索尼A1:奥运年旗舰微单的速度与画质对决,谁能代表微单技术的最高峰?

    佳能EOS R1凭借AI驱动的智能对焦、20张预连拍、机内神经网络降噪和6K RAW视频,结合深度学习技术与专业生态整合,在体育与新闻摄影领域展现出更前瞻的技术高度。 在专业体育与新闻摄影领域,佳能EOS R1和索尼A1是两款代表品牌顶尖技术的旗舰微单。它们都在追求速度、对焦与画质的极致平衡,但实现…

    2026年9月21日
    100
  • laravel如何进行安全的SQL查询以防止注入_Laravel安全SQL查询防注入方法

    使用Eloquent和Query Builder并配合参数绑定可有效防止SQL注入。Laravel通过PDO预处理机制自动转义参数,确保安全;应避免拼接用户输入,尤其在whereRaw等原生语句中需使用?占位符绑定变量;所有用户输入均需验证,对ID类字段强制类型转换,并禁止将用户输入直接用于表名、字…

    2026年9月21日
    000
  • 在Java中如何分析异常堆栈性能开销

    异常堆栈在高并发场景下开销显著,因JVM需遍历调用栈、创建对象、字符串拼接及同步操作,频繁使用将增加GC压力与CPU消耗;可通过JMH测试量化影响,发现填充堆栈耗时可达清空的10倍以上;建议避免在热点代码抛异常、禁用非必要堆栈填充、按需打印日志、使用异步日志框架,并借助JFR、Profiler和GC…

    2026年9月21日
    000
  • PHP/MySQL:高效合并订单商品并按日期分组显示

    本教程将指导如何在PHP/MySQL应用中,将同一日期的订单商品合并显示在同一行,以提高数据展示的清晰度。核心解决方案是利用MySQL的GROUP_CONCAT函数在数据库层面进行高效聚合,避免复杂的PHP逻辑处理,从而简化代码并优化性能。 订单数据展示的常见挑战 在开发在线购物平台时,通常需要向用…

    2026年9月21日
    100
  • google浏览器CPU占用率过高怎么解决_google浏览器CPU占用过高解决方法

    Chrome CPU占用过高可通过清除缓存、禁用高耗能扩展、结束高占用进程、更新浏览器、关闭硬件加速及禁用Software Reporter Tool解决。 如果您在使用Google Chrome浏览器时发现电脑运行缓慢或风扇狂转,很可能是由于Chrome的CPU占用率过高导致系统资源被大量消耗。以…

    2026年9月21日
    000
  • win11任务管理器打不开怎么办_win11任务管理器无法打开修复方法

    1、使用SFC和DISM命令修复系统文件后重启;2、通过gpedit.msc检查并禁用“删除任务管理器”策略;3、在注册表中将DisableTaskMgr值设为0;4、创建新用户账户测试是否解决任务管理器无法打开问题。 如果您尝试打开任务管理器时没有响应或无法启动,可能是由于系统文件损坏、组策略设置…

    2026年9月21日
    000
  • VSCode报错怎么显示中文_VSCode错误信息本地化与中文显示教程

    安装中文语言包可将VSCode界面和错误提示转为中文,提升使用便捷性;但外部工具如编译器、解释器生成的报错仍为英文,因VSCode仅显示其原始输出,无法翻译。 在VSCode中让报错信息显示中文,核心在于安装并启用官方的中文(简体)语言包。这不仅仅是针对错误信息,而是将整个VSCode的用户界面本地…

    2026年9月21日
    000
  • 如何在MindSpore中训练AI大模型?华为AI框架的训练教程

    如何在MindSpore中训练AI大模型?华为AI框架的训练教程如何在MindSpore中训练AI大模型?华为AI框架的训练教程如何在MindSpore中训练AI大模型?华为AI框架的训练教程如何在MindSpore中训练AI大模型?华为AI框架的训练教程

    答案:MindSpore通过自动并行、混合精度、优化器状态分片等技术,结合Profiler工具调试性能瓶颈,实现大模型高效分布式训练。 ☞☞☞AI 智能聊天, 问答助手, AI 智能搜索, 免费无限量使用 DeepSeek R1 模型☜☜☜ 在MindSpore中训练AI大模型,核心在于巧妙地利用其…

    2026年9月21日 • 用户投稿
    300
  • Java ConcurrentSkipListMap在并发场景下应用

    ConcurrentSkipListMap是基于跳跃表实现的线程安全有序映射,支持高并发读写与高效范围查询,适用于需排序的并发场景,如排行榜系统;相比ConcurrentHashMap,它提供有序性与导航操作,但插入查找为O(log n),内存开销较大,适合读多写少或需区间扫描的业务。 在高并发场景…

    2026年9月21日
    100
  • MySQL数据库如何支持多租户业务_设计策略与实现?

    MySQL数据库如何支持多租户业务_设计策略与实现?MySQL数据库如何支持多租户业务_设计策略与实现?MySQL数据库如何支持多租户业务_设计策略与实现?MySQL数据库如何支持多租户业务_设计策略与实现?

    mysql 支持多租户架构的关键在于选择合适的数据隔离策略,并兼顾性能与运维管理。1. 常见方式包括共享数据库共享表(资源利用率高但隔离性差)、共享数据库独立表(平衡隔离性与维护成本)和独立数据库(隔离性强但管理复杂)。2. 租户识别需在请求前确定租户id,并自动附加到sql查询中,可通过视图或中间…

    2026年9月21日 • 用户投稿
    000
  • VSCode怎么启动Layui项目_VSCode运行Layui前端框架项目教程

    必须使用本地服务器运行Layui项目,因为直接打开HTML文件通过file://协议会受浏览器安全限制,导致AJAX、跨域等功能异常,Layui组件无法正常加载;推荐安装Node.js后使用npm全局安装http-server,通过命令行启动服务,或在VSCode中安装Live Server插件,右…

    2026年9月21日
    000
  • 俄罗斯Яндекс账号登录入口 Yandex电脑版官方网站登录

    答案是https://www.yandex.com/。该网站提供搜索、地图、新闻、翻译等服务,界面简洁,支持个性化设置与账户同步,并拥有邮箱、云存储及丰富的应用生态。 1、立即进入“☞☞☞☞点击俄罗斯yandex搜索引擎入口☜☜☜☜”; 2、立即进入“☞☞☞☞点击快速获取Yandex免登录官网链接☜…

    2026年9月21日
    000
  • 蝴蝶号直播掉帧、断流怎么办?技术实用建议

    蝴蝶号直播掉帧、断流怎么办?技术实用建议蝴蝶号直播掉帧、断流怎么办?技术实用建议蝴蝶号直播掉帧、断流怎么办?技术实用建议蝴蝶号直播掉帧、断流怎么办?技术实用建议

    解决蝴蝶号直播掉帧、断流问题需从硬件、软件、网络三方面入手。1. 硬件方面:检查cpu和gpu压力,必要时升级硬件或降低分辨率、帧率;确保摄像头、采集卡、内存正常工作。2. 软件方面:调整分辨率、帧率、码率至合适水平;使用h.265或硬件编码减轻cpu负担;设置关键帧间隔为2秒;关闭后台程序并检查平…

    2026年9月21日 • 用户投稿
    300
  • 如何使用Ribbet的AI功能裁剪图片?快速实现精准图像裁剪

    答案:Ribbet的AI裁剪功能可快速智能识别主体并推荐裁剪方案,支持手动微调与多种比例选择,结合亮度、色彩等编辑工具优化效果,适用于制作符合社交媒体尺寸要求的封面图,操作简便且大部分功能免费,适合追求效率的普通用户。 ☞☞☞AI 智能聊天, 问答助手, AI 智能搜索, 免费无限量使用 DeepS…

    2026年9月21日
    400
  • SystemTap

    SystemTap 简介 systemtap 是一款用于诊断 linux 系统性能或功能问题的开源工具。它使得对运行中的 linux 系统进行诊断和调试变得更加便捷和高效。有了 systemtap,开发者和调试人员无需重新编译内核、安装新内核或重启系统等繁琐步骤。为了解决系统问题或提升性能,开发者只…

    2026年9月21日
    100
  • 卢伟冰:功能手机、智能手机之后 手机行业正进入新周期

    9月4日,小米集团总裁卢伟冰表示,继功能机时代与智能机时代之后,全球手机产业正迈入一个全新时代。 卢伟冰今日在社交平台发文提到:“我从2002年进入手机行业,有幸完整见证了功能手机和智能手机两大发展阶段。如今,AI时代已经到来,整个行业正在酝酿深刻变革,步入全新的发展周期。” 回望过去,功能手机时期…

    2026年9月21日
    200
  • 谷歌浏览器官方主站入口 最新Chrome在线登录页面

    谷歌浏览器官方主站入口是https://www.google.com,该页面具备界面简洁、操作流畅、集成化服务入口和个性化推荐等特点,支持多设备访问且无广告干扰。 谷歌浏览器官方主站入口在哪里?这是不少网友都关注的,接下来由PHP小编为大家带来谷歌浏览器最新Chrome在线登录页面相关信息,感兴趣的…

    2026年9月21日
    000
  • win11怎么退回win10系统_win11降级回win10系统操作教程

    可在10天内通过系统恢复功能退回Windows 10,保留文件但卸载新增应用;超期则需用媒体工具或第三方软件重装,后者操作更简便但会清除数据。 如果您最近将系统升级到 Windows 11,但发现使用不习惯或存在兼容性问题,则可以考虑退回至 Windows 10。在特定时间窗口内,Windows 提…

    2026年9月21日
    000

发表回复

登录后才能评论
关注微信