ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

modern-go/concurrent 源码解析:从 concurrent.Map 到可取消的 Executor 并发模型

modern-go/concurrent 源码解析:从 concurrent.Map 到可取消的 Executor 并发模型 机器学习深度学习数据可视化可观测性【免费下载链接】wandbThe AI developer platform. Use Weights Biases to train and fine-tune models, and manage models from experimentation to production.项目地址https://gitcode.com/gh_mirrors/wa/wandb点击查看免费下载modern-go/concurrent是一个轻量级的 Go 并发工具库以两个核心 API 著称concurrent.Mapsync.Map 的向后移植版本解决 Go 1.9 以下版本的可移植性问题与concurrent.Executor为 goroutine 提供显式所有权、统一取消与 panic 兜底机制。本指南将以本仓库中 vendored 的源码core/vendor/github.com/modern-go/concurrent为蓝本逐一拆解两个 API 的使用方式、底层实现与设计动机帮助你掌握如何让 goroutine 变得可控、可取消、可追踪的工程实践。一、库概览两个 API 解决什么问题该库定位极小只提供两个核心抽象API解决的问题核心能力concurrent.MapGo 1.9 之前没有sync.Map并发安全的 map 写法因版本而异屏蔽版本差异提供与sync.Map一致的使用体验concurrent.Executor裸go关键字启动的 goroutine 无法统一取消、无法统一处理 panic显式所有权 统一取消 panic 回调兜底此外库内还暴露了两个可配置的全局日志器ErrorLogger默认输出到 stderr与InfoLogger默认丢弃用于 executor 运行期间的错误与信息上报见 log.go。二、concurrent.Mapsync.Map 的向后移植1. 为什么需要它sync.Map在 Go 1.9 中才被引入。若项目需要兼容 Go 1.9 以下的编译器直接使用sync.Map会导致编译失败。concurrent.Map通过构建标签build tags在编译期自动选择正确的实现使代码一次编写、处处可编译。2. 双实现机制构建标签分流同一目录下存在两个互为镜像的实现文件通过// build标签隔离Go 1.9 及以上go_above_19.goMap直接内嵌标准库的sync.Map零成本复用官方实现// build go1.9 type Map struct { sync.Map } func NewMap() *Map { return Map{} }Go 1.9 以下go_below_19.go以sync.RWMutex 普通 map 自行实现线程安全读写读操作加读锁、写操作加写锁语义与sync.Map对齐// build !go1.9 type Map struct { lock sync.RWMutex data map[interface{}]interface{} } func NewMap() *Map { return Map{ data: make(map[interface{}]interface{}, 32), } } func (m *Map) Load(key interface{}) (elem interface{}, found bool) { m.lock.RLock() elem, found m.data[key] m.lock.RUnlock() return } func (m *Map) Store(key interface{}, elem interface{}) { m.lock.Lock() m.data[key] elem m.lock.Unlock() }注意两个细节一是低版本实现为 map 预分配了 32 个桶位的初始容量降低扩容频率二是Store在锁内直接覆盖旧值Load返回(elem, found)双返回值与sync.Map的调用约定完全一致。3. 使用方法继承自 README无论最终编译到哪个 Go 版本业务代码的写法完全一致m : concurrent.NewMap() m.Store(hello, world) elem, found : m.Load(hello) // elem 为 world // found 为 true从源码结构看Map通过类型内嵌sync.MapGo ≥ 1.9继承了LoadOrStore、Delete、Range等其余方法在低版本分支中则仅实现了Load/Store两个最小方法因此跨版本代码建议只依赖这两个方法以保证可移植性边界清晰。三、concurrent.Executor显式所有权与可取消的 goroutine1. 设计动机go关键字的缺陷裸go func()启动的 goroutine 存在三个工程痛点无法从外部统一取消panic 一旦逃逸会直接崩溃整个进程goroutine 与启动者之间没有归属关系难以跟踪生命周期。Executor正是针对这三点设计的替代品。接口定义非常克制见 executor.gotype Executor interface { // Go starts a new goroutine controlled by the context Go(handler func(ctx context.Context)) }接口只暴露Go方法停止、等待等管理能力属于具体实现类型如UnboundedExecutor文档注释也明确提示启动并拥有 executor 的一方应当使用具体类型而非接口。2. UnboundedExecutor不限制 goroutine 数量的实现UnboundedExecutorunbounded_executor.go是对活跃 goroutine 数量不加上限的实现其内部状态包括一个context.Contextcancel函数作为统一取消信号源一个map[string]int 互斥锁记录每个 goroutine 的启动位置文件:行号与存活计数用于等待全部退出与问题排查可覆盖的HandlePanic回调默认行为是打日志而非崩溃。它还提供了一个进程级单例GlobalUnboundedExecutor生命周期与程序本身对齐适合承载希望在main退出前统一回收的后台任务。3. 官方 README 示例完整继承executor : concurrent.NewUnboundedExecutor() executor.Go(func(ctx context.Context) { everyMillisecond : time.NewTicker(time.Millisecond) for { select { case -ctx.Done(): fmt.Println(goroutine exited) return case -everyMillisecond.C: // do something } } }) time.Sleep(time.Second) executor.StopAndWaitForever() fmt.Println(executor stopped)把 goroutine 挂到 executor 实例上之后我们可以获得两种能力README 原意统一取消通过Stop/StopAndWait/StopAndWaitForever停止 executor从而取消其名下所有 goroutinepanic 兜底goroutine 内发生的 panic 由回调统一处理默认行为是记录日志不再导致整个应用崩溃。四、源码级机制拆解1. Go启动、追踪、兜底三合一Go方法unbounded_executor.go执行三个关键步骤记录启动位置通过reflect.ValueOf(handler).Pointer()runtime.FuncForPC解析出函数名、文件与行号并在activeGoroutines计数中 1启动带 recover 的 goroutine真正执行的handler(executor.ctx)外层包裹defer recover()任何 panic 都会被捕获panic 分流处理恢复后若recovered ! nil优先调用实例级executor.HandlePanic若被显式覆盖否则回落到包级全局HandlePanic随后将存活计数 -1。包级默认的HandlePanic实现unbounded_executor.go会向ErrorLogger输出 panic 信息与完整堆栈var HandlePanic func(recovered interface{}, funcName string) { ErrorLogger.Println(fmt.Sprintf(%s panic: %v, funcName, recovered)) ErrorLogger.Println(string(debug.Stack())) }2. 三种停止语义方法行为Stop()调用cancel()发出取消信号立即返回不等待goroutine 退出StopAndWait(ctx)取消后循环轮询每 100ms 一次直到所有活跃 goroutine 计数归零若传入的 ctx 先被取消则提前返回避免无限阻塞StopAndWaitForever()等价于StopAndWait(context.Background())即等到全部退出为止轮询过程由checkNoActiveGoroutinesunbounded_executor.go驱动遍历activeGoroutines计数表只要存在计数大于 0 的条目就向InfoLogger打印仍等待 goroutine 退出及对应的启动位置并返回 false 继续等待。这里有一个值得注意的协作约定executor 只负责发出取消信号context 取消并不强杀 goroutine。任务必须自己监听ctx.Done()并返回否则StopAndWaitForever可能无限等待——README 示例中select分支的case -ctx.Done()正是这一契约的标准实现。3. 全局单例与退出纪律GlobalUnboundedExecutor的注释明确告诫它的生命周期与程序对齐期望main函数显式调用 stop它并不能神奇地感知主函数退出。换言之使用全局 executor 的进程必须在退出前调用停止方法否则后台 goroutine 会随进程一起被硬终止丢失优雅回收的机会。五、实战组合一个完整的可运行示例将 Map 与 Executor 组合使用可以构造一个并发写入、统一回收的完整场景package main import ( context fmt time github.com/modern-go/concurrent ) func main() { // 1. 线程安全 map多个 goroutine 并发写入 m : concurrent.NewMap() // 2. 无上限 executor统一管理所有后台任务 executor : concurrent.NewUnboundedExecutor() for i : 0; i 10; i { key : fmt.Sprintf(worker-%d, i) executor.Go(func(ctx context.Context) { ticker : time.NewTicker(time.Millisecond) defer ticker.Stop() for { select { case -ctx.Done(): m.Store(key, stopped) return case -ticker.C: m.Store(key, time.Now().String()) } } }) } time.Sleep(50 * time.Millisecond) executor.StopAndWaitForever() // 取消所有任务并等待全部退出 // 3. 安全读取结果 if v, ok : m.Load(worker-3); ok { fmt.Println(worker-3 , v) } }运行方式在包含本仓库core模块的环境下进入 core 目录执行go run即可验证。该示例完整覆盖了 README 中提到的三个核心诉求——线程安全共享状态Map、统一取消StopAndWaitForever、以及任务退出后的状态收敛。六、使用建议与边界说明确认协程监听取消信号Stop只发信号不回收任务不响应ctx.Done()时StopAndWait*会一直等待设计后台任务时应始终在select中监听ctx.Done()。覆盖 HandlePanic 以自定义告警库提供包级concurrent.HandlePanic与实例级字段两种覆盖入口生产环境可在此接入指标上报或告警系统替换默认的日志输出。低版本 Map 仅保证最小方法集跨 Go 版本代码请只依赖Load/Store如需Delete、Range等能力需自行确认目标版本分支的实现。全局单例须显式停止GlobalUnboundedExecutor不会自动感知main退出进程退出前记得调用其停止方法。日志默认静默InfoLogger默认写入ioutil.Discard见 log.go排查为何一直等待退出这类问题时可将InfoLogger重定向到 stderr 以观察checkNoActiveGoroutines的等待日志。七、在本仓库中的定位本仓库中该库以 vendored 依赖形式存在于 core/vendor/github.com/modern-go/concurrent并在 core/go.mod 中声明为github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd // indirect即 wandb-core 通过modern-go生态引入的传递性依赖其Map与Executor抽象为上层并发代码提供了可移植、可取消的并发原语。对于想要理解 Go 并发封装惯用法的读者这份源码是一份极佳的微型范本构建标签实现多版本兼容、context 驱动的统一取消、recover 兜底与生命周期计数三个模式都浓缩在不足两百行的代码中。赞分享机器学习深度学习数据可视化可观测性【免费下载链接】wandbThe AI developer platform. Use Weights Biases to train and fine-tune models, and manage models from experimentation to production.项目地址https://gitcode.com/gh_mirrors/wa/wandb点击查看免费下载相关推荐KubeSphere 依赖解析modern-go/concurrent 的并发 Map 与可取消 Goroutine Executor 实战指南KubeSphere 依赖解析modern go/concurrent 的并发 Map 与可取消 Goroutine Executor 实战指南 导读 git后端云原生容器编排微服务autoscaler 中 modern-go/concurrent 并发工具库解析concurrent.Map 与 UnboundedExecutor 的源码级指南autoscaler 中 modern go/concurrent 并发工具库解析concurrent.Map 与 UnboundedExecutor 的源码弹性伸缩云原生容器编排深入解析 kops 内置的 modern-go/concurrent跨版本并发 Map 与可取消 Executor 实战指南深入解析 kops 内置的 modern go/concurrent跨版本并发 Map 与可取消 Executor 实战指南 导读 github.com/mo云原生集群管理运维IaC创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表