Go 装饰器模式与洋葱模型

前言 装饰器模式 是实现功能解耦和动态扩展的核心设计模式,而 洋葱模型 作为装饰器模式的链式进阶,是 Web 框架中间件、RPC 拦截器的底层核心。 本文介绍用 Go 实现支持 context.Context 上下文传递的装饰器模式(包括函数式和接口式两种实现),并进阶实现洋葱模型 一、装饰器模式 模式定义 装饰器模式是一种结构型设计模式,核心原则:不修改原有代码逻辑,动态为对象/函数添加额外功能。 它遵循开放封闭原则:对扩展开放,对修改关闭。 核心特性 动态增强:运行期为目标添加前置/后置逻辑 无侵入性:不改动核心业务代码 可组合:多个装饰器链式叠加,自由组合功能 上下文透传:配合 context.Context 实现全链路数据传递(超时、请求ID、用户信息等) 适用场景 日志打印、耗时统计 权限校验、参数校验 缓存、重试、限流熔断 全链路上下文管理 二、Go 实现装饰器模式(带 Context 上下文) 函数式装饰器 Go 没有类继承,通过函数包装实现装饰器是最简洁的方案。 为装饰器加入 context.Context,实现上下文全链路传递(核心需求)。 设计思路 定义统一的业务函数类型(强制携带 Context) 编写核心业务函数(无任何增强逻辑) 编写装饰器函数:接收原函数 → 返回包装后的新函数 装饰器内部通过 Context 传递/读取数据,实现全链路交互 函数式装饰器代码示例 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 package main import ( "context" "fmt" "time" ) // 定义业务函数类型:携带 context.Context,支持上下文传递 // 这是装饰器的统一接口 type BusinessFunc func(ctx context.Context, name string) string // 原始核心业务函数:仅实现核心逻辑,无增强代码 func SayHello(ctx context.Context, name string) string { // 从上下文读取中间件传递的数据 reqID, _ := ctx.Value("req_id").(string) user, _ := ctx.Value("user").(string) fmt.Printf("[核心业务] 请求ID:%s | 操作用户:%s | 执行核心逻辑\n", reqID, user) return fmt.Sprintf("你好,%s!", name) } // 装饰器1:日志装饰器(携带上下文) func LogDecorator(f BusinessFunc) BusinessFunc { return func(ctx context.Context, name string) string { // 前置增强:打印请求日志 reqID := ctx.Value("req_id").(string) fmt.Printf("[日志装饰器] 请求ID:%s | 开始调用,参数:%s\n", reqID, name) // 调用原函数(透传上下文) result := f(ctx, name) // 后置增强:打印响应日志 fmt.Printf("[日志装饰器] 请求ID:%s | 调用结束,结果:%s\n", reqID, result) return result } } // 装饰器2:计时装饰器(携带上下文) func TimeDecorator(f BusinessFunc) BusinessFunc { return func(ctx context.Context, name string) string { start := time.Now() reqID := ctx.Value("req_id").(string) // 调用原函数(透传上下文) result := f(ctx, name) // 后置增强:打印耗时 fmt.Printf("[计时装饰器] 请求ID:%s | 执行耗时:%s\n", reqID, time.Since(start)) return result } } // 装饰器3:上下文初始化装饰器(往ctx存入数据,供下游使用) func ContextDecorator(f BusinessFunc) BusinessFunc { return func(ctx context.Context, name string) string { // 给上下文添加请求ID、用户信息(全链路透传) ctx = context.WithValue(ctx, "req_id", "REQ_123456") ctx = context.WithValue(ctx, "user", "admin") fmt.Println("[上下文装饰器] 初始化上下文完成") return f(ctx, name) // 无后置增强 } } func main() { // 链式装饰:顺序 = 上下文装饰 → 日志装饰 → 计时装饰 decoratedFunc := ContextDecorator(LogDecorator(TimeDecorator(SayHello))) // 根上下文 ctx := context.Background() // 执行增强后的函数 res := decoratedFunc(ctx, "张三") fmt.Println("\n最终返回结果:", res) } 运行结果 注意每层装饰器中被装饰函数 BusinessFunc 的调用位置,后面的代码是被装饰函数返回后再执行的 ...

August 18, 2022 · 7 min · 1283 words · erpan

golang context包用法理解

同时启很多个goroutine来完成一个任务,在一些必要的情况下如何跟踪或取消这些goroutine?常见的有下面几种方式: WaitGroup,goroutine之间的同步 for select 加上 stop channel来监听消息管理协程 context包 context包就是用来在goroutine之间传递上下文信息的,它提供了超时timeout和取消cancel机制,利用了channel进行信息传递,方便的管理具有复杂层级关系的多个goroutine,这些复杂的层级关系类似于一棵树,可以有多个分叉。Go标准库中的net/http,database/sql等都用到了context包。 接口 context包定义了两个接口 Context 1 2 3 4 5 6 type Context interface { Deadline() (deadline time.Time, ok bool) Done() <-chan struct{} Err() error Value(key interface{}) interface{} } Deadline 方法,第一个返回值是该上下文执行完的截止时间,第二个返回值是布尔值,表示是否设置了截止日期,没设置返回ok==false,否则为true。该方法幂等 Done,返回一个只读通道,在context被取消或者context到了deadline时,此通道会被关闭 WithCancel 上下文,当调用cancel后此通道关闭关闭 WithDeadline,到达deadline时间点后自动关闭此通道 WithTimeout,指定的时间超时后自动关闭此通道 Err,Done返回的通道未被关闭时,Err的返回值是nil,关闭后,返回 context 取消的原因,被取消时返回Canceled,到了截止时间返回 DeadlineExceeded Value,获取该Context上绑定的值,需要注意的是键和值都是any类型。以上四个方法都是幂等的 canceler ,私有接口,表示一个可以被取消cancel的context对象,包内的*cancelCtx 和 *timerCtx结构体实现了此接口 1 2 3 4 type canceler interface { cancel(removeFromParent bool, err, cause error) Done() <-chan struct{} } 创建 几个实现此接口的结构体关系: ...

February 15, 2022 · 3 min · 503 words · erpan

ArgoCD试用

ArgoCD 遵循GitOPS模式,应用定义、配置和环境等都应该是声明式和版本化的。应用部署和生命周期管理是自动化、可审计和易于理解的 特性 拥有GitOPS的一切特性,如回滚到任意git commit点 自动发布应用到指定环境,支持多集群管理 支持多种配置管理工具、模板 (Kustomize,Helm, Jsonnet, plain-YAML,自定义配置管理插件) 支持单点登录 (OIDC, OAuth2, LDAP, SAML 2.0, GitHub, GitLab, Microsoft, LinkedIn) 支持多租户及RBAC授权 服务健康状态分析 资源版本偏移检查和Web UI实时可视化 可自动或手动同步资源到目标状态 提供命令行工具,自动化集成简单方便 Webhook 集成 (GitHub, BitBucket, GitLab) 支持访问令牌 支持各阶段钩子定义PreSync, Sync, PostSync hooks to support complex application rollouts (e.g.blue/green & canary upgrades) 应用事件和API调用审计追踪 有暴露Prometheus 指标 Parameter overrides for overriding helm parameters in Git 几个核心概念 Application:指定manifest路径下的一组kubernetes资源,是一个CRD AppProject:Application的逻辑分组,可以配置一些约束选项,如限制git源、目标集群、namespace和资源类型等,也可以订阅project roles App of Apps:官网文档中的 cluster bootstrapping,一个Application包含多个子Application,在批量创建Application时比较有用,也可以用来自我管理(官方文档有示例) 下图是定义了一个Application对象,其指定的仓库目录下包含多种资源,其中的子Application对象所指定的仓库目录又可以包含多种资源,我们可以开启递归的选项,在往仓库添加资源对象的时候自动apply到对应集群当中去 注意点: 一个ArgoCD实例中,Application名字是唯一的,且只能放在与argo部署的同一名称空间中 Application中没有指定resources-finalizer.argocd.argoproj.io终结器,在删除Application的时候是不会删除它所管理的资源,App of Apps也是一样 ApplicationSet 跨集群和仓库灵活的管理Applications,补充了以集群管理为中心的场景 ...

January 10, 2022 · 3 min · 485 words · erpan

client-go SharedInformer

在client-go中,informer是一个核心的概念,用于从Kubernetes API服务器中读取对象。它可以监视一个或多个API对象,并在对象发生变化时自动更新本地缓存。 Informer主要组成部分: sharedIndexInformer sharedIndexInformer是informer的核心组件。它负责从API服务器获取对象并更新本地缓存。在更新缓存后,它会触发事件处理器,以便通知其他组件对象已经更新。 informerSyncHandler informerSyncHandler是sharedIndexInformer的事件处理程序,它将缓存中的对象与API服务器中的对象进行比较,并更新缓存。当处理完所有的更新操作后,informerSyncHandler会将更新后的对象发送到事件队列。 informerEventHandler informerEventHandler是事件处理程序的接口,用于接收事件并执行特定的操作。当sharedIndexInformer接收到更新事件时,它将调用informerEventHandler来处理事件 indexer indexer是informer的本地缓存,用于存储从API服务器获取的对象。indexer使用map数据结构存储对象,其中键是对象的名称,值是对象本身。此外,indexer还使用索引数据结构(例如Set、List)来优化查找和筛选操作。每个indexer都关联一个ObjectStore,ObjectStore是一个更高级别的抽象,用于管理一组API对象 watcher watcher是informer的事件源,用于从API服务器中获取对象并将它们发送到事件队列中。watcher实现了Kubernetes API服务器上的watch机制,它会定期向API服务器发送请求,以获取当前对象的状态。在获取状态后,watcher将对比上一次请求的结果,找到任何发生变化的对象,并将它们发送到事件队列中 在运行时,informer将indexer和watcher组合在一起,以便实现对象的监视和更新。当watcher从API服务器中获取到对象时,它会将它们添加到indexer中。如果对象已经存在于indexer中,则watcher将更新对象的状态。一旦对象被更新,informerSyncHandler将被触发,以便将更新后的对象发送到事件队列中。在事件处理程序中,可以通过调用indexer对象的方法来访问缓存中的对象 画了一个可能不太标准的图: 了解了Informer架构和Reflector、DeltaFIFO、Indexer几个组件后,再回头来看示例中的SharedInformer使用,它如何将前面几个组件关联起来 示例程序 创建Informer对象 注册事件处理程序 启动Informer 以下面代码为例: 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 func main() { var err error var config *rest.Config var kubeconfig *string if home := homedir.HomeDir(); home != "" { kubeconfig = flag.String("kubeconfig", filepath.Join(home, ".kube", "config"), "[可选] kubeconfig 绝对路径") } else { kubeconfig = flag.String("kubeconfig", "", "kubeconfig 绝对路径") } flag.Parse() // 初始化 *rest.Config 对象 if config, err = rest.InClusterConfig(); err != nil { if config, err = clientcmd.BuildConfigFromFlags("", *kubeconfig); err != nil { panic(err.Error()) } } // 创建 *ClientSet 对象 clientSet, err := kubernetes.NewForConfig(config) if err != nil { panic(err.Error()) } // 初始化 SharedInformerFactory,暂且10s relist一次 sharedInformerFactory := informers.NewSharedInformerFactory(clientSet, time.Second*10) // 每种kubernetes内置资源都实现了Informer // 如创建podInformer:podInformer := sharedInformerFactory.Core().V1().Pods() // 下面创建deploymentInformer deployInformer := sharedInformerFactory.Apps().V1().Deployments() // 实际上调用了 InformerFor(&appsv1.Deployment{}, f.defaultInformer) 构造 deployInformer 中的 SharedIndexInformer informer := deployInformer.Informer() // 创建 DeploymentLister,拥有list方法 //type deploymentLister struct { // indexer cache.Indexer //} deployLister := deployInformer.Lister() // 注册事件处理程序 informer.AddEventHandler(cache.ResourceEventHandlerFuncs{ AddFunc: onAdd, UpdateFunc: onUpdate, DeleteFunc: onDelete, }) // 程序进程退出前,通知informer退出,informer以goroutine运行 stopper := make(chan struct{}) defer close(stopper) // 启动 informer 去 list & watch // 即 SharedIndexInformer 的 run方法,new了DeltaFIFO,然后new出sharedIndexInformer.controller,调用Controller的run方法(即Reflector的run方法),run里面创建了Reflector并执行 sharedInformerFactory.Start(stopper) // 等待所有启动的 Informer 的缓存被同步 sharedInformerFactory.WaitForCacheSync(stopper) informer.Run(stopper) // 从本地缓存中获取 default 中的所有 deployment 列表 deployments, err := deployLister.Deployments("default").List(labels.Everything()) if err != nil { panic(err) } for idx, deploy := range deployments { fmt.Printf("%d -> %s\n", idx+1, deploy.Name) } <-stopper } func onAdd(obj interface{}) { deploy := obj.(*v1.Deployment) fmt.Println("add a deployment:", deploy.Name) } func onUpdate(old, new interface{}) { oldDeploy := old.(*v1.Deployment) newDeploy := new.(*v1.Deployment) fmt.Println("update deployment:", oldDeploy.Name, newDeploy.Name) } func onDelete(obj interface{}) { deploy := obj.(*v1.Deployment) fmt.Println("delete a deployment:", deploy.Name) } 执行逻辑 NewSharedInformerFactory创建Informer对象,函数中传入了clientset和defaultResync参数,clientset是访问APIServer的具体实现,resyncPeriod指定了informer在从API服务器中获取数据之间等待的时间,sharedInformerFactory的informers字段是一个map结构,每种资源类型为key,对应value是SharedIndexInformer 1 2 3 4 5 6 7 8 9 10 func NewSharedInformerFactory(client kubernetes.Interface, defaultResync time.Duration) SharedInformerFactory factory := &sharedInformerFactory{ client: client, namespace: v1.NamespaceAll, defaultResync: defaultResync, // 示例中我们传入的是10s informers: make(map[reflect.Type]cache.SharedIndexInformer), startedInformers: make(map[reflect.Type]bool), customResync: make(map[reflect.Type]time.Duration), } 每种kubernetes内置资源都已经实现了资源对应的Informer,通过sharedInformerFactory.Apps().V1().Deployments()链式调用构造出deployInformer,遵循组、版本、资源(即GVK)格式,其他k8s资源的informer创建也是一样的操作 1 2 3 4 5 6 type deploymentInformer struct { factory internalinterfaces.SharedInformerFactory // TweakListOptionsFunc is a function that transforms a v1.ListOptions tweakListOptions internalinterfaces.TweakListOptionsFunc namespace string } deployInformer.Informer() 用deployInformer创建deployment对应的SharedIndexInformer。这里面初始化了重要的两个对象:cache.Indexers和ListerWatcher 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 // 给InformerFor传入两个参数: // 1. Deployment对象空结构体 // 2. Deployment的默认Informer构造函数(它在调用后返回SharedIndexInformer) // 返回Deployment对应的SharedIndexInformer func (f *deploymentInformer) Informer() cache.SharedIndexInformer { return f.factory.InformerFor(&appsv1.Deployment{}, f.defaultInformer) } // InformerFor的参数defaultInformer方法的实现,返回了构造出的SharedIndexInformer func (f *deploymentInformer) defaultInformer(client kubernetes.Interface, resyncPeriod time.Duration) cache.SharedIndexInformer { // 传入了参数cache.Indexers // Indexers的键是索引分类名,为常量namespace,值为IndexFunc,取出对象的namespace值 return NewFilteredDeploymentInformer(client, f.namespace, resyncPeriod, cache.Indexers{cache.NamespaceIndex: cache.MetaNamespaceIndexFunc}, f.tweakListOptions) } func NewFilteredDeploymentInformer(client kubernetes.Interface, namespace string, resyncPeriod time.Duration, indexers cache.Indexers, tweakListOptions internalinterfaces.TweakListOptionsFunc) cache.SharedIndexInformer { // 传入了ListerWatcher对象和Indexer参数来构造SharedIndexInformer return cache.NewSharedIndexInformer( &cache.ListWatch{ // ListFunc/WatchFunc函数就是来list和watch APIServer资源的动作,具体实现在clientset里面 ListFunc: func(options metav1.ListOptions) (runtime.Object, error) { if tweakListOptions != nil { tweakListOptions(&options) } return client.AppsV1().Deployments(namespace).List(context.TODO(), options) }, WatchFunc: func(options metav1.ListOptions) (watch.Interface, error) { if tweakListOptions != nil { tweakListOptions(&options) } return client.AppsV1().Deployments(namespace).Watch(context.TODO(), options) }, }, &appsv1.Deployment{}, resyncPeriod, // NewSharedInformerFactory时传入的10s indexers, ) } func NewSharedIndexInformer(lw ListerWatcher, exampleObject runtime.Object, defaultEventHandlerResyncPeriod time.Duration, indexers Indexers) SharedIndexInformer { realClock := &clock.RealClock{} // 最终返回的informer是这个sharedIndexInformer sharedIndexInformer := &sharedIndexInformer{ processor: &sharedProcessor{clock: realClock}, indexer: NewIndexer(DeletionHandlingMetaNamespaceKeyFunc, indexers), listerWatcher: lw, objectType: exampleObject, resyncCheckPeriod: defaultEventHandlerResyncPeriod, // 10s defaultEventHandlerResyncPeriod: defaultEventHandlerResyncPeriod, // 10s cacheMutationDetector: NewCacheMutationDetector(fmt.Sprintf("%T", exampleObject)), clock: realClock, } return sharedIndexInformer } // 入参: // 1. 通用对象Object,此处即Deployment{} // 2. deployment资源类型对应的sharedIndexInformer的构造函数 func (f *sharedInformerFactory) InformerFor(obj runtime.Object, newFunc internalinterfaces.NewInformerFunc) cache.SharedIndexInformer { f.lock.Lock() defer f.lock.Unlock() // 反射取得对象类型为Deployment informerType := reflect.TypeOf(obj) informer, exists := f.informers[informerType] if exists { return informer } resyncPeriod, exists := f.customResync[informerType] if !exists { resyncPeriod = f.defaultResync // 默认是零值,所以赋予我们传入的10s } // 如果没有此类型的informer的话就新建一个,此例是deployment类型的informer informer = newFunc(f.client, resyncPeriod) // 填充了sharedInformerFactory的informers字段 // informers map[reflect.Type]cache.SharedIndexInformer f.informers[informerType] = informer // 返回sharedIndexInformer实例 return informer } 创建 DeploymentLister,它里面就一个字段indexer,indexer是索引器对象,可以在本地缓存中根据索引取数。如从本地缓存indexer中获取 default 名称空间的所有 deployment 列表:deployments, err := deployLister.Deployments("default").List(labels.Everything()) 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 func (f *deploymentInformer) Lister() v1.DeploymentLister { // 传入一个参数为cache.Indexer,即前面func (f *deploymentInformer) defaultInformer方法中 // 初始化的cache.Indexers{cache.NamespaceIndex: cache.MetaNamespaceIndexFunc} return v1.NewDeploymentLister(f.Informer().GetIndexer()) } func (s *sharedIndexInformer) GetIndexer() Indexer { return s.indexer } func NewDeploymentLister(indexer cache.Indexer) DeploymentLister { return &deploymentLister{indexer: indexer} } // deploymentLister implements the DeploymentLister interface. type deploymentLister struct { indexer cache.Indexer } 至此,我们有了sharedIndexInformer,其已拥有重要的字段cache.Indexers和ListerWatcher已初始化,但目前还没有启动ListerWatcher去从APIServer拿数据,没有把这些数据包装为Delta缓存下来,没有存入到Indexer,也没有后续对数据的处理逻辑 ...

November 15, 2021 · 10 min · 2098 words · erpan

client-go RingGrowingBuffer 环形缓冲区

在client-go源码中,processorListener对象里面定义了一个RingBuffer用于缓存所有尚未分发的事件通知,在此记录下这个RingBuffer。 RingBuffer一般用于数据的缓存机制,例如tcp协议里面数据包的缓冲就利用到了RingBuffer。 client-go中的这个buffer是非线程安全、可增长、无边界的先进先出环形缓冲区。环是一个逻辑上的概念,有了环,此段内存空间就可以重复利用,不用频繁重新申请内存,本质上数据还是存在数组里面的,这个数组的大小可以按需进行倍数扩容,扩容后需要重新分配内存空间并拷贝未消费的数据到新数组来。因为它是数组,内存是预先分配的,数组是内存上连续的一段空间,它有一个容易预测的访问模式,因此对CPU高速缓存友好,垃圾回收(GC)在这种情况下也不用做什么。 在实现上,可以理解为两个指针:a) 读指针、b)写指针。在一段buffer上,读指针控制下一次该读数据的位置,写指针控制下一次该写数据的位置。在数组里面我们可以直接用数组下标。要遵循FIFO原则,读指针不能超过写指针,两指针重叠了要么buffer写满了,要么buffer为空。 参考这里的一张图 目前我们只需要考虑: 存储啥数据类型 数据如何存放 何时空间满了需要扩容 不能丢失未消费数据 k8s中的源码: 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 // 源码路径k8s.io/utils/buffer/ring_growing.go package buffer // 非线程安全、可增长的环形缓冲区 type RingGrowing struct { data []interface{} // 存任意数据的数组 n int // 缓冲区大小 beg int // 第一个可用的元素位置索引 readable int // 未消费的元素数量 } // 初始化一个RingBuffer func NewRingGrowing(initialSize int) *RingGrowing { return &RingGrowing{ data: make([]interface{}, initialSize), n: initialSize, } } // 取出未消费元素中的第一个元素 func (r *RingGrowing) ReadOne() (data interface{}, ok bool) { if r.readable == 0 { return nil, false } r.readable-- element := r.data[r.beg] r.data[r.beg] = nil // Remove reference to the object to help GC if r.beg == r.n-1 { // 这种情况就是读到数组最后一个元素了,下次读就得从列表的头部开始,以免越界 r.beg = 0 } else { r.beg++ } return element, true } // 在buffer尾部添加一个元素,buffer满了就扩容 func (r *RingGrowing) WriteOne(data interface{}) { // 满了的情况 if r.readable == r.n { newN := r.n * 2 // 新开辟一个两倍大小的数组 newData := make([]interface{}, newN) to := r.beg + r.readable if to <= r.n { copy(newData, r.data[r.beg:to]) } else { // 未消费的元素在数组两端的情况 copied := copy(newData, r.data[r.beg:]) copy(newData[copied:], r.data[:(to%r.n)]) } r.beg = 0 r.data = newData r.n = newN } r.data[(r.readable+r.beg)%r.n] = data r.readable++ } client-go中的这个buffer是非线程安全、可增长、无边界的先进先出环形缓冲区,因此,在此client-go ringBuffer基础上可以继续考虑下面几个问题: ...

October 11, 2021 · 2 min · 252 words · erpan