并发模式

简介: 并发模式并不是一种函数的运用、亦或者实际存在的东西。他是前人对于并发场景的运用总结与经验。他与23中设计模式一样。好啦,话不多说。开干

并发模式并不是一种函数的运用、亦或者实际存在的东西。他是前人对于并发场景的运用总结与经验。他与23中设计模式一样。好啦,话不多说。开干


无论是如何厉害的架构还是编程方式,我始终相信都是从零开始,不断的抽象,不断的迭代的。抽象思维对于我们尤为重要。那么我们也带着这样的一个疑问。思考到底什么是抽象


首先我们将要学习的是work pool模式


work pool


不知道大家是否在go并发的时候遇见过以下几个问题或者想法


  • goroutine的数量控制可能并不是那么称心如意


  • goroutine,创造过多,造成资源浪费。且并发效果也并非那么好。他正如正态分布那样。到达某个极点所带来的收益将会下降


  • goroutine复用的问题,往往一个goroutine都只处理了一个任务。不断的创建与删除


  • 甚至更多。。。


workpool,首先分析以上问题,我个人总结都以上其实是一个问题,groutine与任务死死的绑定,并没有进行解耦。比如像这样。


// example
package main
import (
    "fmt"
    "time"
)
func exs(accept <-chan int, recipient chan<- int) {
    for result := range accept {
        fmt.Println("Received only sent channel a:", result)
        recipient <- result + 2
    }
    //fmt.Println("Send Only", recipient)
}
func main() {
    startTime := time.Now()
    ch := make(chan int, 10)
    for i := 0; i < 100; i++ {
        go func(ch <-chan int) {
            time.Sleep(time.Second * 5)
            fmt.Println(<-ch)
        }(ch)
        ch <- i
    }


那么我们来改造一下,然后进行代码剖析。代码如下


package main
import (
    "fmt"
    "time"
)
func work(id int, jobs <-chan int, result chan<- int) {
    for j := range jobs {
        fmt.Println("Worker [ID]", id, "Start Process JoB [Id]", j)
        time.Sleep(time.Second * 2)
        //fmt.Println("Working, will Spend 2 s")
        fmt.Println("Worker [ID]", id, "Carry Process JoB [Id]", j)
        result <- j * 2
    }
}
func main() {
    const jobNumber = 1000
    const workerNumber = 100
    jobs := make(chan int, workerNumber)
    result := make(chan int, jobNumber)
    // Create Worker(start Goroutines)
    for w := 0; w <= workerNumber; w++ {
        go work(w, jobs, result)
    }
    // arrange work
    for j := 0; j <= jobNumber; j ++ {
        jobs <- j
    }
    // 获取结果
    for r := 0; r <= jobNumber; r ++ {
        <- result
    }
}


work pool的精髓在于将任务,与groutine进行分离。只关心初始的任务与结果。是不是与函数式编程很像呢?我也这么觉得,嘻嘻


来吧,我们剖析一下代码


  1. 首先我们定义了两个常量(建议是常量),jobNumworkerNumber,故名思义他们分别是任务数量,以及工人数量。你可以将他们看出生产者与消费者。


  1. 我们定义了两个channel,他们作为我们发送指令与获取结果的通道。记得加缓存哦,否则将造成死锁


  1. 最后就是分别定义消费者-groutine,生产者jobNumber,然后传递任务进入goroutine。然后我们就只需要得到结果就好啦


nice,虽然很简单。但也有无限的可能性哦。你还可以进一步抽象,变成一个通用的goroutine pool。


Pipeline 模式


Pipeline 模式也称为流水线模式,模拟的就是现实世界中的流水线生产。


从技术上看,每一道工序的输出,就是下一道工序的输入,在工序之间传递的东西就是数据,这种模式称为流水线模式,而传递的数据称为数据流。下面我们用代码模拟柴火烧饭的过程


package main
import "fmt"
func main() {
    combust := wash(10)
    rice := combustion(combust)
    packs := open(rice)
    //输出测试,看看效果
    for p := range packs {
        fmt.Println(p)
    }
}
func wash(n int) <-chan string {
    out := make(chan string)
    go func() {
        defer close(out)
        for i := 1; i <= n; i++ {
            out <- fmt.Sprint("洗米", i)
        }
    }()
    return out
}
func combustion(in <-chan string) <-chan string {
    out := make(chan string)
    go func() {
        defer close(out)
        for c := range in {
            out <- "烧饭(" + c + ")"
        }
    }()
    return out
}
func open(in <-chan string) <-chan string {
    out := make(chan string)
    go func() {
        defer close(out)
        for c := range in {
            out <- "开锅(" + c + ")"
        }
    }()
    return out
}


开锅(烧饭(洗米1))

开锅(烧饭(洗米2))

开锅(烧饭(洗米3))

开锅(烧饭(洗米4))

开锅(烧饭(洗米5))

开锅(烧饭(洗米6))

开锅(烧饭(洗米7))

开锅(烧饭(洗米8))

开锅(烧饭(洗米9))

开锅(烧饭(洗米10))


首先,我为什么一定强调是柴火烧饭呢,难道柴火香一点?那可不,必须的。


其实这里,我们需要思考一个问题,什么是可异步的,什么是不可异步的?


拓展:


可异步:例如网络请求,发送网络请求后,立马发送下一个。尽量减少网络io阻塞,从而提高效率。可前提是,网络io阻塞可以不用等待


不可异步:也就是说我们每一步都必须参与其中,计算机它无法独自去完成。例如柴火烧饭,没柴火咋烧饭,魔法么。当然你硬要说火烧一次就一直可以不需要人去干预,那咱也没办法了不是


在这里,生产者与消费者可能并不像之前那么分的那么开了,首先


洗米(生产者)


烧饭(消费者、生产者)


开锅(消费者)


这种模式称为流水线模式,而传递的数据称为数据流


分治模式


就像前面所说那样,每一道必须依靠前面完成了才能进行下一步,但我们发现其中烧饭或者太慢了,我们可以分而治之,然后合并。也可以达到我们需要的效果。


package main
import (
    "fmt"
    "sync"
    "time"
)
func main() {
    combust := wash(10)
    rice1 := combustion(combust)
    rice2 := combustion(combust)
    rice3 := combustion(combust)
    rice := merge(rice1, rice2, rice3)
    packs := open(rice)
    //输出测试,看看效果
    for p := range packs {
        fmt.Println(p)
    }
}
func wash(n int) <-chan string {
    out := make(chan string)
    go func() {
        defer close(out)
        for i := 1; i <= n; i++ {
            out <- fmt.Sprint("洗米", i)
        }
    }()
    return out
}
func combustion(in <-chan string) <-chan string {
    out := make(chan string)
    go func() {
        defer close(out)
        time.Sleep(2)
        for c := range in {
            out <- "烧饭(" + c + ")"
        }
    }()
    return out
}
func open(in <-chan string) <-chan string {
    out := make(chan string)
    go func() {
        defer close(out)
        for c := range in {
            out <- "开锅(" + c + ")"
        }
    }()
    return out
}
func merge(ins ...<-chan string) <-chan string {
    var wg sync.WaitGroup
    out := make(chan string)
    //把一个channel中的数据发送到out中
    p := func(in <-chan string) {
        defer wg.Done()
        for c := range in {
            out <- c
        }
    }
    wg.Add(len(ins))
    //扇入,需要启动多个goroutine用于处于多个channel中的数据
    for _, cs := range ins {
        go p(cs)
    }
    //等待所有输入的数据ins处理完,再关闭输出out
    go func() {
        wg.Wait()
        close(out)
    }()
    return out
}


Futures 模式


Pipeline 流水线模式中的工序是相互依赖的,上一道工序做完,下一道工序才能开始。但是在我们的实际需求中,也有大量的任务之间相互独立、没有依赖,所以为了提高性能,这些独立的任务就可以并发执行。


举个例子,比如我打算自己做顿火锅吃,那么就需要洗菜、烧水。洗菜、烧水这两个步骤相互之间没有依赖关系,是独立的,那么就可以同时做,但是最后做火锅这个步骤就需要洗好菜、烧好水之后才能进行。这个做火锅的场景就适用 Futures 模式。


Futures 模式可以理解为未来模式,主协程不用等待子协程返回的结果,可以先去做其他事情,等未来需要子协程结果的时候再来取,如果子协程还没有返回结果,就一直等待


Futures 模式下的协程和普通协程最大的区别是可以返回结果,而这个结果会在未来的某个时间点使用。所以在未来获取这个结果的操作必须是一个阻塞的操作,要一直等到获取结果为止。


如果你的大任务可以拆解为一个个独立并发执行的小任务,并且可以通过这些小任务的结果得出最终大任务的结果,就可以使用 Futures 模式。


Referer


22讲通关go语言-飞雪无情

目录
相关文章
|
存储 缓存 NoSQL
【分布式】Redis与Memcache的对比分析
【1月更文挑战第25天】【分布式】Redis与Memcache的对比分析
|
7月前
|
安全 应用服务中间件 Linux
HTTPS 优化完整方案解析
本文详解HTTPS性能优化全方案,从原理到实操,涵盖硬件加速(AES-NI)、软件升级(内核与OpenSSL)及协议层优化(TLS 1.3、ECDSA、会话复用等),配合Nginx配置模板与验证方法,助你实现安全与速度双提升,显著降低访问延迟。
1523 156
|
6月前
|
人工智能 前端开发 搜索推荐
AI 英语口语 APP 的费用
AI英语口语APP开发费用分三块:研发人力(60%-70%)、AI模型调用(Token计费)、第三方授权。预算分三级:MVP版15–30万(基础对话纠错);进阶版40–100万(数字人+发音打分);企业版150万+(自研模型+VR沉浸)。2026年可借轻量化模型、端侧VAD、跨平台框架降本。#AI教育 #AI英语
|
7月前
|
数据采集 人工智能 搜索推荐
智能体来了:降本增效的终极杀手锏,销售人必看的生存指南
内容摘要:AI智能体(AI Agents)正重塑销售底层逻辑。本文深度拆解智能体如何通过自动化获客、个性化触达及全天候线索转化,协助销售人员打破内卷僵局,实现业绩呈指数级增长。
565 3
|
7月前
|
存储 人工智能 弹性计算
2026年阿里云建站费用与功能全指南:自助建站、模板建站、开发型建站方案详解
在数字化需求日益增长的当下,搭建网站成为个人展示、企业推广的重要途径。阿里云针对不同技术基础与业务规模,推出 “自购服务器建站”“万小智 AI 模板建站”“云企业官网定制建站” 三种核心方案,价格从 38 元 / 年到数万元 / 年不等,覆盖从个人到中大型企业的全场景需求。本文基于最新官方定价与实测数据,从方案细节、价格体系、功能对比、场景适配等维度展开解析,为用户提供客观选型参考。
|
9月前
|
人工智能 JavaScript IDE
别用"战术勤奋"掩盖"战略懒惰":AI时代的降维竞品分析
5%的产品死于"盲视"。本文不仅是一套竞品分析AI指令,更是一次从战术勤奋到战略觉醒的认知升级。教你如何利用AI构建全天候商业情报雷达,寻找巨头缝隙中的差异化生存之道,实现商业战场的降维打击。
792 7
|
11月前
|
机器学习/深度学习 编解码 Python
Python图片上采样工具 - RealESRGANer
Real-ESRGAN基于深度学习实现图像超分辨率放大,有效改善传统PIL缩放的模糊问题。支持多种模型版本,推荐使用魔搭社区提供的预训练模型,适用于将小图高质量放大至大图,放大倍率越低效果越佳。
819 3
|
11月前
|
数据采集 JSON 数据挖掘
淘宝API对接系列:商品详情与评论数据分析(JSON数据返回)
1. 商品详情API(taobao.item.get) • 功能:获取商品基础信息(标题、价格、库存、销量)、图片、类目、促销信息等。
|
11月前
|
Java 数据库连接 开发者
Spring Framework 核心技术详解
本文档旨在深入解析 Java Spring Framework 的核心技术原理与应用。与侧重于快速开发的 Spring Boot 不同,本文将聚焦于 Spring 框架本身的设计理念、核心容器、控制反转(IoC)、面向切面编程(AOP)、数据访问与事务管理等基础且强大的模块。通过理解这些核心概念,开发者能够更深刻地领悟 Spring 生态系统的设计哲学,并具备解决复杂企业级应用开发问题的能力。
597 4
|
12月前
|
人工智能 搜索推荐 数据挖掘
小红书电商 API 开启小红书店铺电商内容营销新范式
小红书电商 API 为商家提供自动化运营与内容营销新范式,支持商品管理、批量发布 UGC、数据分析等功能,提升效率、降低成本,助力品牌实现精准营销与用户深度互动。