Go 互斥锁 Mutex 源码分析(一)

前言 锁作为并发编程中的关键一环,是应该要深入掌握的。 锁 示例 实现锁很简单,示例如下: 1var global int 2 3func main() { 4 var mu sync.Mutex 5 var wg sync.WaitGroup 6 7 for i := 0; i < 2; i++ { 8 wg.Add(1) 9 go func(i int) { 10 defer wg.Done() 11 mu.Lock() 12 global++ 13 mu.Unlock() 14 }(i) 15 } 16 17 wg.Wait() 18 fmt.Println(global) 19} 输出: ...

2024-08-23 · xhy

client-go DeltaFIFO 精讲

0. 介绍 client-go list&watch 精讲 介绍过 Reflector 通过 list&watch 实时获取 Kubernetes 资源的信息。并且作为生产者将资源信息存到 DeltaFIFO 中。 那么消费者是在哪里定义的呢?本文围绕 DeltaFIFO 继续介绍是哪个模块消费了 DeltaFIFO 中的资源。 ...

2024-08-21 · xhy

client-go event handler 精讲

1. 介绍 client-go DeltaFIFO 介绍了从 DeltaFIFO 中取出的资源将被放到 indexer ,接着交由 processorListener 处理。 本文继续看 processorListener 是如何处理 DeltaFIFO pop 的资源的。 2. 启动 handler processorListener 的启动在 sharedIndexInformer.Run() 方法定义: func (s *sharedIndexInformer) Run(stopCh <-chan struct{}) { ... // 开启协程运行 sharedIndexInformer.processor.run 方法 wg.StartWithChannel(processorStopCh, s.processor.run) ... } func (p *sharedProcessor) run(stopCh <-chan struct{}) { func() { p.listenersLock.RLock() defer p.listenersLock.RUnlock() for listener := range p.listeners { // 开启协程运行 processorListener.run p.wg.Start(listener.run) // 开启协程运行 processorListener.pop p.wg.Start(listener.pop) } p.listenersStarted = true }() <-stopCh ... } sharedProcessor.run 的重点在 processorListener.pop 和 processorListener.run,分别介绍如下。 ...

2024-08-21 · xhy

client-go indexer 精讲

1. 介绍 indexer 是 client-go 中的缓存,不同于普通缓存,indexer 是带索引的缓存。 在 client-go DeltaFIFO 精讲 中介绍到 DeltaFIFO 中的元素被加入到 threadSafeMap 中,这个 map 实际是 indexer 的一部分。 indexer 在 client-go 定义为接口,实现该 indexer 接口的是 cache 对象: ...

2024-08-21 · xhy

client-go list & watch 精讲

1. 介绍 client-go 架构 通过 client-go 可以实现访问 kubernetes 资源的实时性,可靠性和顺序性。那么 client-go 是 如何实现的呢? 2. list & watch client-go 通过 list&watch 机制实现实时访问 kubernetes 资源,满足实时性。接下来从源码层面解析 client-go 的 list& watch 机制。 ...

2024-08-21 · xhy

client-go workqueue 精讲

1. 介绍 传递到 event handler 的资源一般不会立即处理,而是先放到 workqueue 中。这是因为业务的处理速度要比资源的更新速度慢。 通过 workqueue 可解耦业务的处理逻辑和资源更新逻辑。 ...

2024-08-21 · xhy

containerd 源码分析:创建 container(三)

文接 containerd 源码分析:创建 container(二) 启动 task 上节介绍了创建 task,task 创建之后将返回 response 给 ctr。接着,ctr 调用 task.Start 启动容器。 1// containerd/client/task.go 2func (t *task) Start(ctx context.Context) error { 3 r, err := t.client.TaskService().Start(ctx, &tasks.StartRequest{ 4 ContainerID: t.id, 5 }) 6 if err != nil { 7 ... 8 } 9 t.pid = r.Pid 10 return nil 11} 12 13// containerd/api/services/tasks/v1/tasks_grpc.pb.go 14func (c *tasksClient) Start(ctx context.Context, in *StartRequest, opts ...grpc.CallOption) (*StartResponse, error) { 15 out := new(StartResponse) 16 err := c.cc.Invoke(ctx, "/containerd.services.tasks.v1.Tasks/Start", in, out, opts...) 17 if err != nil { 18 return nil, err 19 } 20 return out, nil 21} ctr 调用 contaienrd 的 /containerd.services.tasks.v1.Tasks/Start 接口创建 task。进入 containerd 查看提供该服务的插件: ...

2024-06-04 · xhy

containerd 源码分析:创建 container(二)

文接 containerd 源码分析:创建 container(一) 创建容器进程 创建 container 成功后,接着创建 task, task 将根据 container metadata 创建容器进程。 创建 task 进入 tasks.Newtask 创建 task: 1// containerd/cmd/ctr/commands/tasks/tasks_unix.go 2func NewTask(ctx gocontext.Context, client *containerd.Client, container containerd.Container, checkpoint string, con console.Console, nullIO bool, logURI string, ioOpts []cio.Opt, opts ...containerd.NewTaskOpts) (containerd.Task, error) { 3 ... 4 t, err := container.NewTask(ctx, ioCreator, opts...) 5 if err != nil { 6 return nil, err 7 } 8 ... 9} 10 11// containerd/client/container.go 12func (c *container) NewTask(ctx context.Context, ioCreate cio.Creator, opts ...NewTaskOpts) (_ Task, err error) { 13 ... 14 t := &task{ 15 client: c.client, 16 io: i, 17 id: c.id, 18 c: c, 19 } 20 ... 21 response, err := c.client.TaskService().Create(ctx, request) 22 if err != nil { 23 return nil, errdefs.FromGRPC(err) 24 } 25 t.pid = response.Pid 26 return t, nil 27} 类似创建 container,这里调用 container.client.TaskService().Create 创建 task: ...

2024-06-04 · xhy