Golang管道模式实现 channel流水线案例

使用Go的channel和goroutine实现整数处理流水线:生成整数→平方→过滤大于100→输出结果,各阶段通过channel串联,每个阶段函数接收输入channel并返回输出channel,形成可组合、高效的数据处理链。

golang管道模式实现 channel流水线案例

在Go语言中,channel和goroutine的组合非常适合实现“流水线(pipeline)”模式。流水线将数据处理过程拆分为多个阶段,每个阶段由一个或多个goroutine通过channel连接,形成一条高效、可扩展的数据处理链。下面通过一个具体的案例来展示如何用Golang实现channel流水线。

1. 场景说明:整数处理流水线

假设我们需要实现一个数据处理流程:生成一组整数 → 平方每个数 → 过滤出大于100的结果 → 输出最终结果。我们可以将每个步骤封装为一个独立的函数,通过channel串联起来。

2. 阶段函数实现

每个阶段接收一个输入channel,返回一个输出channel,形成可组合的流水线结构。

立即学习“go语言免费学习笔记(深入)”;

阶段1:生成数据

gen函数生成一组整数并通过channel输出。

func gen(nums …int)   out := make(chan int, len(nums))
  go func() {
    for _, n := range nums {
      out     }
    close(out)
  }()
  return out
}

立即学习“go语言免费学习笔记(深入)”;

阶段2:平方处理

square函数从输入channel读取数据,计算平方后发送到输出channel。

func square(in   out := make(chan int)
  go func() {
    for n := range in {
      out     }
    close(out)
  }()
  return out
}

立即学习“go语言免费学习笔记(深入)”;

阶段3:过滤大于100的数据

filterOver100函数只传递大于100的值。

func filterOver100(in   out := make(chan int)
  go func() {
    for n := range in {
      if n > 100 {
        out       }
    }
    close(out)
  }()
  return out
}

3. 组合流水线并运行

将各个阶段通过channel串联起来,形成完整的处理流程。

func main() {
  // 构建流水线
  source := gen(1, 2, 3, 4, 5, 6, 7, 8, 9, 10) // 1~10
  squared := square(source) // 平方
  filtered := filterOver100(squared) // 过滤 >100
  
  // 消费结果
  for result := range filtered {
    fmt.Println(result)
  }
}

输出结果:
121 (11²)
144 (12²)
169 (13²)
196 (14²)
225 (15²)
… 以此类推,实际输入是1~10,平方后最大为100,因此实际无输出。若想看到输出,可将gen改为 gen(11,12,13) 或调整输入。

4. 优化:支持多阶段并行与扇出/扇入

在高并发场景中,可以对某个阶段启动多个worker,提高处理能力。

func squareParallel(in   out := make(chan int, workers)
  
  var wg sync.WaitGroup
  for i := 0; i     wg.Add(1)
    go func() {
      for n := range in {
        out       }
      wg.Done()
    }()
  }
  
  go func() {
    wg.Wait()
    close(out)
  }()
  return out
}

这种模式称为“扇出(fan-out)”和“扇入(fan-in)”,可以显著提升处理吞吐量。

基本上就这些。Golang的channel流水线模式简洁而强大,适合ETL、数据清洗、消息处理等场景。关键是每个阶段职责单一,通过channel自然解耦,易于测试和扩展。不复杂但容易忽略的是资源清理和goroutine泄漏问题,确保所有channel最终被关闭,避免阻塞。

以上就是Golang管道模式实现 channel流水线案例的详细内容,更多请关注创想鸟其它相关文章!

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

(0)
打赏 微信扫一扫 微信扫一扫 支付宝扫一扫 支付宝扫一扫
如何用Golang反射实现依赖注入 构建简易IoC容器实战
上一篇 2025年12月15日 16:21:47
Golang服务降级方案 优雅应对高负载
下一篇 2025年12月15日 16:22:07

相关推荐

  • 怎么在VSCode里配置Go语言环境?

    安装Go并配置环境变量后,在VSCode中安装官方Go扩展,通过命令面板安装gopls、delve等必要工具,并设置保存时自动格式化与导入,即可实现代码补全、格式化和调试功能。 在 VSCode 中配置 Go 语言开发环境其实不复杂,只要安装好工具链并正确设置,就能获得代码补全、格式化、调试等完整功…

    2026年9月20日
    100
  • 如何为VSCode配置Go语言开发环境?

    首先安装Go环境并验证版本与环境变量,然后在VSCode中安装官方Go插件,接着通过命令行手动安装gopls和dlv等关键工具,最后创建测试文件确认语法高亮、代码补全和调试功能正常即可完成配置。 为 VSCode 配置 Go 语言开发环境其实不难,只要正确安装工具和插件,就能获得代码补全、跳转、格式…

    2026年9月12日
    100
  • Workerman如何实现消息队列?WorkermanRabbitMQ集成?

    Workerman通过与RabbitMQ集成,利用其常驻内存和事件驱动特性,实现高效的消息生产与消费。相比传统PHP-FPM每次请求重建连接,Workerman在onWorkerStart中建立持久连接,复用连接资源,显著降低开销,提升吞吐量和实时性。作为消费者,Workerman可实时监听队列,消…

    2026年9月11日
    100
  • 游戏数据分析:PHP+Go组合如何高效处理海量打点数据?

    高效游戏数据分析:PHP和Go的完美结合 一款游戏数据分析系统的设计中,开发者选择了PHP和Go语言的组合方案。PHP负责后台分析系统,而Go语言则承担打点接口和数据处理的重任。 挑战:海量并发打点数据的处理 游戏运行过程中,大量的并发打点操作会产生海量数据。为了应对这一挑战,开发者计划利用Kafk…

    2026年9月1日
    100
  • PHP+Go游戏打点分析系统如何优化性能?

    提升PHP和Go游戏数据分析系统性能的策略 本文探讨如何优化一个由PHP后端分析系统、Go语言打点接口、Kafka异步计算以及MySQL数据库组成的游戏数据分析系统。该系统的设计逻辑清晰,但性能方面存在改进空间。 避免直接数据库写入:性能瓶颈的突破 当前架构中,Go打点接口直接写入MySQL数据库,…

    2026年9月1日
    200
  • 高并发游戏打点分析:PHP+Go组合如何高效处理海量数据?

    高效游戏打点分析:PHP和Go的完美结合 本文探讨如何构建一个高效的游戏打点分析系统,以应对高并发和海量数据带来的挑战。我们将重点介绍一种基于PHP和Go的组合方案,并分析其优缺点及改进建议。 系统架构: 本系统采用PHP和Go协同工作,数据处理流程如下: 立即学习“PHP免费学习笔记(深入)”; …

    2026年9月1日
    200
  • Docker:应用容器引擎 Docker简介,Docker安装与启动(一步一步教你安装,不相信有看了这个教程还不会的人)

    一、%ignore_a_1%简介 1.1 什么是Docker Docker 是一个用Go语言开发的开源容器项目。通过利用操作系统现有的机制和特性,它实现了比传统虚拟机更轻量级的虚拟化(简单来说,Docker内嵌一个极小的系统,例如Linux仅需5M左右,Windows亦如此)。Docker实现的是内…

    2026年8月28日
    100
  • 协程栈(Coroutine Stack)的内存管理

    协程栈的内存管理是通过用户态栈和运行时环境来实现的。1)在python中,协程使用生成器和yield机制,共享全局解释器锁,需处理暂停和恢复逻辑。2)在go中,goroutine使用m:n调度模型,运行时自动调整栈大小,防止栈溢出和内存泄漏。 在编程世界中,协程栈(Coroutine Stack)的…

    2026年8月28日
    100
  • 分布式运维监控系统 WGCLOUD v3.3.6 全新发布 详细解读更新功能点

    分布式运维监控系统 WGCLOUD v3.3.6 全新发布 详细解读更新功能点分布式运维监控系统 WGCLOUD v3.3.6 全新发布 详细解读更新功能点分布式运维监控系统 WGCLOUD v3.3.6 全新发布 详细解读更新功能点分布式运维监控系统 WGCLOUD v3.3.6 全新发布 详细解读更新功能点

    wgcloud是一款功能强大且易于使用的分布式运维监控系统,具有易部署、轻量级和高效的特点。其server端基于springboot开发,而agent端则采用go语言编写。该系统的核心功能包括:监控主机系统信息、cpu使用率、cpu温度、内存使用情况、网络流量、磁盘i/o、磁盘空间、系统负载、硬盘s…

    2026年8月27日 用户投稿
    100
  • golang怎么连接mysql数据库

    golang操作mysql 安装 go get “github.com/go-sql-driver/mysql”go get “github.com/jmoiron/sqlx” 连接数据库 var Db *sqlx.DBdb, err := sqlx.Open(“mysql”,”username:p…

    用户投稿 2026年8月26日
    000
  • 如何解决HEIC/AVIF图片转换难题?使用Composer和heif-converter轻松搞定!

    可以通过一下地址学习composer:学习地址 告别 HEIC/AVIF 图片兼容性烦恼:用 Composer 玩转 heif-converter 相信很多朋友都有过这样的经历:朋友用 iphone 拍了张照片发给你,结果你发现它是个 .heic 文件。或者,你在网上下载了一些高质量的图片,发现它们…

    用户投稿 2026年8月26日
    100
  • Workerman的未来路线图

    workerman未来将专注于提升性能、扩展多语言支持、加强生态系统集成和提高易用性。1.通过优化底层实现和网络协议提升性能。2.逐步支持go、python等语言。3.加强与docker、kubernetes的集成。4.推出更多工具和文档提高易用性。 关于Workerman的未来路线图,我认为Wor…

    2026年8月25日
    000
  • Go Template中实现异步表单提交:避免页面刷新

    本文将指导如何在Go模板中实现异步表单提交,以避免传统表单提交导致的页面整体刷新。通过利用JavaScript的`FormData`对象结合AJAX技术(如Axios或原生Fetch API),用户可以提交表单数据而无需重新加载整个页面,从而显著提升用户体验和应用的响应速度。 异步表单提交原理与实践…

    2025年12月23日
    100
  • Go模板中实现表单异步提交与页面无刷新技术指南

    本教程详细介绍了如何在%ignore_a_1%模板中实现表单的异步提交,避免页面整体刷新。通过利用javascript的`event.preventdefault()`阻止默认提交行为,结合`formdata`对象收集表单数据,并使用`axios`或`fetch`等http客户端库发送异步请求,从而…

    2025年12月23日
    000
  • 利用Ajax在Go模板中实现表单无刷新提交

    本文详细介绍了如何在go模板中实现表单的异步提交,从而避免页面整体重载。通过利用javascript的`formdata`对象和`axios`等http客户端,我们可以拦截表单的默认提交行为,将数据以异步请求的方式发送到后端,显著提升用户体验和页面响应速度。 引言:提升Go模板表单交互体验 在Web…

    2025年12月23日
    000
  • Go模板中实现表单无刷新提交:利用AJAX优化用户体验

    本文将详细介绍如何在go模板或其他html页面中实现表单的无刷新提交。通过拦截默认的表单提交事件,利用javascript的formdata对象和ajax技术(如axios或fetch),将表单数据异步发送到服务器,从而避免页面整体重载,显著提升用户体验和应用性能。 在传统的Web应用中,当用户提交…

    2025年12月23日
    000
  • HTML5的WebSocket是什么?如何建立实时通信?

    HTML5的WebSocket是什么?如何建立实时通信?HTML5的WebSocket是什么?如何建立实时通信?HTML5的WebSocket是什么?如何建立实时通信?HTML5的WebSocket是什么?如何建立实时通信?

    websocket与传统http请求/长轮询的本质区别在于通信模式和效率。1. 传统http请求是“一问一答”式的单向通信,每次请求都需要重新建立连接,效率低;2. http长轮询虽然延长了等待时间,但本质上仍是请求-响应模型,连接在每次数据传输后断开,依然存在延迟和资源浪费;3. websocke…

    2025年12月22日 用户投稿
    200
  • Node.js与区块链项目中CP-ABE实现策略:跨语言方案与集成考量

    本文探讨了在Node.%ignore_a_1%和区块链项目中实现密文策略属性基加密(CP-ABE)所面临的挑战,指出JavaScript生态中缺乏维护良好的原生库。文章详细介绍了Python、Rust、C++和Go等语言中成熟的CP-ABE库,并提出了跨语言集成策略及在区块链环境中应用CP-ABE的…

    2025年12月21日
    100
  • 在Node.js与区块链项目中实现CP-ABE的策略与方案

    本文探讨了在Node.js和区块链项目中实现密文策略属性基加密(CP-ABE)所面临的库选择挑战。鉴于JavaScript生态中缺乏维护良好的直接CP-ABE库,文章提出了利用Python、Rust、C++或Go等语言中的成熟库,并通过微服务架构进行集成的实用策略,同时提供了概念性代码示例和在区块链…

    2025年12月21日
    000
  • CP-ABE在Node.js与区块链应用中的实现路径探究

    CP-ABE在Node.js和区块链项目中的实现面临JavaScript库稀缺的挑战。本文将探讨当前主流的CP-ABE库生态,指出Python、C++和Rust等语言中的成熟解决方案,并讨论Node.js绑定及Go语言库作为替代方案的可行性,为开发者提供跨语言集成的策略与建议,以克服JavaScri…

    2025年12月21日
    100

发表回复

登录后才能评论
关注微信