Golang源码分析之golang/sync之singleflight

2023-11-05 00:36

本文主要是介绍Golang源码分析之golang/sync之singleflight,希望对大家解决编程问题提供一定的参考价值,需要的开发者们随着小编来一起学习吧!

1.1. 项目介绍

golang/sync库拓展了官方自带的sync库,提供了errgroup、semaphore、singleflight及syncmap四个包,本次分析singlefliht的源代码。
singlefliht用于解决单机协程并发调用下的重复调用问题,常与缓存一起使用,避免缓存击穿。

1.2.使用方法

go get -u golang.org/x/sync

  • 核心API:Do、DoChan、Forget
  • Do:同一时刻对某个Key方法的调用, 只能由一个协程完成,其余协程阻塞直到该协程执行成功后,直接获取其生成的值,以下是一个避免缓存击穿的常见使用方法:
func main() {var flight singleflight.Groupvar errGroup errgroup.Group// 模拟并发获取数据缓存for i := 0; i < 10; i++ {i := ierrGroup.Go(func() error {fmt.Printf("协程%v准备获取缓存\n", i)v, err, shared := flight.Do("getCache", func() (interface{}, error) {// 模拟获取缓存操作fmt.Printf("协程%v正在读数据库获取缓存\n", i)time.Sleep(100 * time.Millisecond)fmt.Printf("协程%v读取数据库生成缓存成功\n", i)return "mockCache", nil})if err != nil {fmt.Printf("err = %v", err)return err}fmt.Printf("协程%v获取缓存成功, v = %v, shared = %v\n", i, v, shared)return nil})}if err := errGroup.Wait(); err != nil {fmt.Printf("errGroup wait err = %v", err)}
}
// 输出:只有0号协程实际生成了缓存,其余协程读取生成的结果
协程0准备获取缓存
协程4准备获取缓存
协程3准备获取缓存
协程2准备获取缓存
协程6准备获取缓存
协程5准备获取缓存
协程7准备获取缓存
协程1准备获取缓存
协程8准备获取缓存
协程9准备获取缓存
协程0正在读数据库获取缓存
协程0读取数据库生成缓存成功
协程0获取缓存成功, v = mockCache, shared = true
协程8获取缓存成功, v = mockCache, shared = true
协程2获取缓存成功, v = mockCache, shared = true
协程6获取缓存成功, v = mockCache, shared = true
协程5获取缓存成功, v = mockCache, shared = true
协程7获取缓存成功, v = mockCache, shared = true
协程9获取缓存成功, v = mockCache, shared = true
协程1获取缓存成功, v = mockCache, shared = true
协程4获取缓存成功, v = mockCache, shared = true
协程3获取缓存成功, v = mockCache, shared = true

DoChan:将执行结果返回到通道中,可通过监听通道结果获取方法执行值,这个方法相较于Do来说的区别是执行DoChan后不会阻塞到其中一个协程完成任务,而是异步执行任务,最后需要结果时直接从通道中获取,避免长时间等待。

func testDoChan() {var flight singleflight.Groupvar errGroup errgroup.Group// 模拟并发获取数据缓存for i :=; i < 10; i++ {i := ierrGroup.Go(func() error {fmt.Printf("协程%v准备获取缓存\n", i)ch := flight.DoChan("getCache", func() (interface{}, error) {// 模拟获取缓存操作fmt.Printf("协程%v正在读数据库获取缓存\n", i)time.Sleep( * time.Millisecond)fmt.Printf("协程%v读取数据库获取缓存成功\n", i)return "mockCache", nil})res := <-chif res.Err != nil {fmt.Printf("err = %v", res.Err)return res.Err}fmt.Printf("协程%v获取缓存成功, v = %v, shared = %v\n", i, res.Val, res.Shared)return nil})}if err := errGroup.Wait(); err != nil {fmt.Printf("errGroup wait err = %v", err)}
}
// 输出结果
协程准备获取缓存
协程准备获取缓存
协程准备获取缓存
协程准备获取缓存
协程准备获取缓存
协程准备获取缓存
协程准备获取缓存
协程准备获取缓存
协程准备获取缓存
协程正在读数据库获取缓存
协程读取数据库获取缓存成功
协程准备获取缓存
协程获取缓存成功, v = mockCache, shared = true
协程获取缓存成功, v = mockCache, shared = true
协程获取缓存成功, v = mockCache, shared = true
协程获取缓存成功, v = mockCache, shared = true
协程获取缓存成功, v = mockCache, shared = true
协程获取缓存成功, v = mockCache, shared = true
协程获取缓存成功, v = mockCache, shared = true
协程获取缓存成功, v = mockCache, shared = true
协程获取缓存成功, v = mockCache, shared = true
协程获取缓存成功, v = mockCache, shared = true

2.源码分析

2.1.项目结构

  • singleflight.go:核心实现,提供相关API
  • singleflight_test.go:相关API单元测试

2.2.数据结构

  • singleflight.go
// singleflight.Group
type Group struct {mu sync.Mutex       // map的锁m  map[string]*call // 保存每个key的调用
}// 一次Do对应的响应结果
type Result struct {Val    interface{}Err    errorShared bool
}// 一个key会对应一个call
type call struct {wg sync.WaitGroupval interface{} // 保存调用的结果err error       // 调用出现的err// 该call被调用的次数dups  int// 每次DoChan时都会追加一个chan在该列表chans []chan<- Result
}

2.3.API代码流程

func (g *Group) Do(key string, fn func() (interface{}, error)) (v interface{}, err error, shared bool)

func (g *Group) Do(key string, fn func() (interface{}, error)) (v interface{}, err error, shared bool) {g.mu.Lock()if g.m == nil {// 第一次执行Do的时候创建mapg.m = make(map[string]*call)}// 已经存在该key,对应后续的并发调用if c, ok := g.m[key]; ok {// 执行次数自增c.dups++g.mu.Unlock()// 等待执行fn的协程完成c.wg.Wait()// ...// 返回执行结果return c.val, c.err, true}// 不存在该key,说明第一次调用,初始化一个callc := new(call)// wg添加,后续其他协程在该wg上阻塞c.wg.Add()// 保存key和call的关系g.m[key] = cg.mu.Unlock()// 真正执行fn函数g.doCall(c, key, fn)return c.val, c.err, c.dups >
}func (g *Group) doCall(c *call, key string, fn func() (interface{}, error)) {normalReturn := falserecovered := false// 第三步、最后的设置和清理工作defer func() {// ...g.mu.Lock()defer g.mu.Unlock()// 执行完成,调用wg.Done,其他协程此时不再阻塞,读到fn执行结果c.wg.Done()// 二次校验map中key的值是否为当前call,并删除该keyif g.m[key] == c {delete(g.m, key)}// ...// 如果c.chans存在,则遍历并写入执行结果for _, ch := range c.chans {ch <- Result{c.val, c.err, c.dups >}}}}()// 第一步、执行fn获取结果func() {//、如果fn执行过程中panic,将c.err设置为PanicErrordefer func() {if !normalReturn {if r := recover(); r != nil {c.err = newPanicError(r)}}}()//、执行fn,获取到执行结果c.val, c.err = fn()//、设置正常返回结果标识normalReturn = true}()// 第二步、fn执行出错,将recovered标识设置为trueif !normalReturn {recovered = true}
}

func (g *Group) DoChan(key string, fn func() (interface{}, error)) <-chan Result

func (g *Group) DoChan(key string, fn func() (interface{}, error)) <-chan Result {// 一次调用对应一个chanch := make(chan Result,)g.mu.Lock()if g.m == nil {// 第一次调用,初始化mapg.m = make(map[string]*call)}// 后续调用,已存在keyif c, ok := g.m[key]; ok {// 调用次数自增c.dups++// 将chan添加到chans列表c.chans = append(c.chans, ch)g.mu.Unlock()// 直接返回chan,不等待fn执行完成return ch}// 第一次调用,初始化call及chans列表c := &call{chans: []chan<- Result{ch}}// wg加一c.wg.Add()// 保存key及call的关系g.m[key] = cg.mu.Unlock()// 异步执行fn函数go g.doCall(c, key, fn)// 直接返回该chanreturn ch
}

3.总结

  • singleflight经常和缓存获取配合使用,可以缓解缓存击穿问题,避免同一时刻单机大量的并发调用获取数据库构建缓存
  • singleflight的实现很精简,核心流程就是使用map保存每次调用的key与call的映射关系,每个call中通过wg控制只存在一个协程执行fn函数,其他协程等待执行完成后,直接获取执行结果,在执行完成后会删去map中的key
  • singleflight的Do方法会阻塞直到fn执行完成,DoChan方法不会阻塞,而是异步执行fn,并通过通道来实现结果的通知

这篇关于Golang源码分析之golang/sync之singleflight的文章就介绍到这儿,希望我们推荐的文章对编程师们有所帮助!



http://www.chinasem.cn/article/346237

相关文章

性能分析之MySQL索引实战案例

文章目录 一、前言二、准备三、MySQL索引优化四、MySQL 索引知识回顾五、总结 一、前言 在上一讲性能工具之 JProfiler 简单登录案例分析实战中已经发现SQL没有建立索引问题,本文将一起从代码层去分析为什么没有建立索引? 开源ERP项目地址:https://gitee.com/jishenghua/JSH_ERP 二、准备 打开IDEA找到登录请求资源路径位置

JAVA智听未来一站式有声阅读平台听书系统小程序源码

智听未来,一站式有声阅读平台听书系统 🌟&nbsp;开篇:遇见未来,从“智听”开始 在这个快节奏的时代,你是否渴望在忙碌的间隙,找到一片属于自己的宁静角落?是否梦想着能随时随地,沉浸在知识的海洋,或是故事的奇幻世界里?今天,就让我带你一起探索“智听未来”——这一站式有声阅读平台听书系统,它正悄悄改变着我们的阅读方式,让未来触手可及! 📚&nbsp;第一站:海量资源,应有尽有 走进“智听

Java ArrayList扩容机制 (源码解读)

结论:初始长度为10,若所需长度小于1.5倍原长度,则按照1.5倍扩容。若不够用则按照所需长度扩容。 一. 明确类内部重要变量含义         1:数组默认长度         2:这是一个共享的空数组实例,用于明确创建长度为0时的ArrayList ,比如通过 new ArrayList<>(0),ArrayList 内部的数组 elementData 会指向这个 EMPTY_EL

如何在Visual Studio中调试.NET源码

今天偶然在看别人代码时,发现在他的代码里使用了Any判断List<T>是否为空。 我一般的做法是先判断是否为null,再判断Count。 看了一下Count的源码如下: 1 [__DynamicallyInvokable]2 public int Count3 {4 [__DynamicallyInvokable]5 get

SWAP作物生长模型安装教程、数据制备、敏感性分析、气候变化影响、R模型敏感性分析与贝叶斯优化、Fortran源代码分析、气候数据降尺度与变化影响分析

查看原文>>>全流程SWAP农业模型数据制备、敏感性分析及气候变化影响实践技术应用 SWAP模型是由荷兰瓦赫宁根大学开发的先进农作物模型,它综合考虑了土壤-水分-大气以及植被间的相互作用;是一种描述作物生长过程的一种机理性作物生长模型。它不但运用Richard方程,使其能够精确的模拟土壤中水分的运动,而且耦合了WOFOST作物模型使作物的生长描述更为科学。 本文让更多的科研人员和农业工作者

MOLE 2.5 分析分子通道和孔隙

软件介绍 生物大分子通道和孔隙在生物学中发挥着重要作用,例如在分子识别和酶底物特异性方面。 我们介绍了一种名为 MOLE 2.5 的高级软件工具,该工具旨在分析分子通道和孔隙。 与其他可用软件工具的基准测试表明,MOLE 2.5 相比更快、更强大、功能更丰富。作为一项新功能,MOLE 2.5 可以估算已识别通道的物理化学性质。 软件下载 https://pan.quark.cn/s/57

工厂ERP管理系统实现源码(JAVA)

工厂进销存管理系统是一个集采购管理、仓库管理、生产管理和销售管理于一体的综合解决方案。该系统旨在帮助企业优化流程、提高效率、降低成本,并实时掌握各环节的运营状况。 在采购管理方面,系统能够处理采购订单、供应商管理和采购入库等流程,确保采购过程的透明和高效。仓库管理方面,实现库存的精准管理,包括入库、出库、盘点等操作,确保库存数据的准确性和实时性。 生产管理模块则涵盖了生产计划制定、物料需求计划、

衡石分析平台使用手册-单机安装及启动

单机安装及启动​ 本文讲述如何在单机环境下进行 HENGSHI SENSE 安装的操作过程。 在安装前请确认网络环境,如果是隔离环境,无法连接互联网时,请先按照 离线环境安装依赖的指导进行依赖包的安装,然后按照本文的指导继续操作。如果网络环境可以连接互联网,请直接按照本文的指导进行安装。 准备工作​ 请参考安装环境文档准备安装环境。 配置用户与安装目录。 在操作前请检查您是否有 sud

线性因子模型 - 独立分量分析(ICA)篇

序言 线性因子模型是数据分析与机器学习中的一类重要模型,它们通过引入潜变量( latent variables \text{latent variables} latent variables)来更好地表征数据。其中,独立分量分析( ICA \text{ICA} ICA)作为线性因子模型的一种,以其独特的视角和广泛的应用领域而备受关注。 ICA \text{ICA} ICA旨在将观察到的复杂信号

Spring 源码解读:自定义实现Bean定义的注册与解析

引言 在Spring框架中,Bean的注册与解析是整个依赖注入流程的核心步骤。通过Bean定义,Spring容器知道如何创建、配置和管理每个Bean实例。本篇文章将通过实现一个简化版的Bean定义注册与解析机制,帮助你理解Spring框架背后的设计逻辑。我们还将对比Spring中的BeanDefinition和BeanDefinitionRegistry,以全面掌握Bean注册和解析的核心原理。