kubelet
# 1.主要功能
# 1.1.Pod管理
kubelet以PodSpec的方式工作,用于描述Pod的YAML或JSON对象。其中,每个Node节点都会启动自己的kubelet进程处理Master节点下发的任务,并采用一组通过各种机制提供的PodSpecs(主要来自api server)确保Pod正常健康运行。 --- 官方提供了几种方式获取容器信息 apiserver: 基于监听etcd目录获取数据 file:启动参数--config指定的配置目录下的文件 http:利用网络从某个地址获取信息 当kubelet监听到etcd中有新的绑定到本节点的Pod时,会按照pod清单的要求创建Pod。同样地,监听到Pod修改时也会根据Pod清单伴随修改Pod状态。1
2
3
4
5
6
7
# 1.2.容器健康检查
--- LivenessProbe 通知kubelet容器当前健康状态,如果LivenessProbe探针探测到容器不健康kubelet将删除该容器并根据重启策略对Pod进行处理 --- ReadinessProbe 判断容器是否启动完成并准备好接受流量,如果ReadinessProbe探针探测到失败额Pod状态会被修改,且Endpoint Controller会从Service的Endpoint中删除容器所在Pod的IP和Endpoint条目1
2
3
4
5
# 1.3.容器监控
kubelet通过cAdvisor获取所在节点及容器数据。cAdvisor是一个开源的分析容器资源使用率和性能特征的代理工具,已经集成到kubelet中,负责监控当前Node节点的信息。cAdvisor自动查找所在节点的容器,自动采集CPU、内存、文件系统和网络的使用统计信息,通过所在节点的Root容器采集并分析节点机的全面使用情况。1
# 2.工作原理
# 2.1.核心机制
kubelet的核心在于循环控制,即SyncLoop。驱动整个控制循环的事件包括:pod更新事件、pod生命周期变化、kubelet本身设置的执行周期、定时清理事件等。 SyncLoop循环会注册很多manager,例如用于探测Pod健康状态的probeManager;用于维护Pod Status的statusManager;用于报告容器创建、失败事件的containerRefManager等。1
2
kubelet调用下层容器运行时的具体过程时,利用CRI(Container Runtime Interface,容器运行时)的gRPC接口实现执行过程管理。 CRI是Kubernetes抽象出的一层接口,旨在解耦容器引擎(不管是docker还是rkt),底层的容器只需要实现CRI接口并暴露gRPC服务即可被kubelet管理。1
2CRI接口目前包含两类,用于镜像增删的ImageService和用于容器管理的RuntimeService

# 2.2.入口
kubelet的入口在cmd/kubelet/kubelet.go,此位置是kubelet命令的server端入口。真正处理逻辑位于pkg\kubelet\kubelet.go下的Run()方法,该方法以协程形式调用。--- 主函数的cli.Run(command)作为处理命令入口 func main() { command := app.NewKubeletCommand() code := cli.Run(command) os.Exit(code) } --- cli.Run(command)通过内部实现的run()处理命令 func Run(cmd *cobra.Command) int { if logsInitialized, err := run(cmd); err != nil { ... } } --- 内部实现的run()执行命令 func run(cmd *cobra.Command) (logsInitialized bool, err error) { ... err = cmd.Execute() return } --- Execute()内部利用ExecuteC()执行命令 func (c *Command) Execute() error { _, err := c.ExecuteC() return err } --- ExecuteC()调用内部实现execute() func (c *Command) ExecuteC() (cmd *Command, err error) { ... err = cmd.execute(flags) ... } --- command.execute()实现上根据传入模块执行不同逻辑,此处以RunE()为例 func (c *Command) execute(a []string) (err error) { ... if c.RunE != nil { if err := c.RunE(c, argWoFlags); err != nil { return err } } else { c.Run(c, argWoFlags) } ... } --- RunE()是cmd/kubelet/app/server.go的方法 RunE: func(cmd *cobra.Command, args []string) error { ... // run the kubelet return Run(ctx, kubeletServer, kubeletDeps, utilfeature.DefaultFeatureGate) }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
注意
kubelet的cmd程序作为Server部分,当集群中的kubelet client创建资源时,由server端获取并交给pkg\kubelet\kubelet.go模块处理
# 3.启动流程
# 3.1.Run
// Run runs the specified KubeletServer with the given Dependencies. func Run(...) error { ... // 启动入口 run(ctx, s, kubeDeps, featureGate) ... return nil }1
2
3
4
5
6
7
8
# 3.2.run
run方法是整个kubelet的核心入口,主要初始化各个重要模块——cadvisor、ContainerManager及RuntimeService,启动kubelet。func run(...) (err error) { ... // 初始化cadvisor监控系统 if kubeDeps.CAdvisorInterface == nil { imageFsInfoProvider := cadvisor.NewImageFsInfoProvider(s.RemoteRuntimeEndpoint) kubeDeps.CAdvisorInterface, err = cadvisor.New(imageFsInfoProvider, s.RootDirectory, cgroupRoots, cadvisor.UsingLegacyCadvisorStats(s.RemoteRuntimeEndpoint)) ... } ... // 初始化containerManager,用于CPU/设备插件/内存分配/资源预留/cgroup等资源管理 if kubeDeps.ContainerManager == nil { ... kubeDeps.ContainerManager = cm.NewContainerManager(...) ... } ... // 初始化容器运行时 kubelet.PreInitRuntimeService(&s.KubeletConfiguration, kubeDeps, s.RemoteRuntimeEndpoint, s.RemoteImageEndpoint) ... // 启动kubelet RunKubelet(s, kubeDeps, s.RunOnce) ... return nil } // PreInitRuntimeService will init runtime service before RunKubelet. func PreInitRuntimeService(...) error { ... // 初始化容器运行时 kubeDeps.RemoteRuntimeService = remote.NewRemoteRuntimeService(remoteRuntimeEndpoint, kubeCfg.RuntimeRequestTimeout.Duration) ... // 初始化镜像运行时 kubeDeps.RemoteImageService = remote.NewRemoteImageService(remoteImageEndpoint, kubeCfg.RuntimeRequestTimeout.Duration) ... return nil }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
# 3.2.RunKubelet
RunKubelet调用createAndInitKubelet初始化kubelet实例及启动GC,最后调用startKubelet启动server及各模块主流程。func RunKubelet(kubeServer *options.KubeletServer, kubeDeps *kubelet.Dependencies, runOnce bool) error { ... // 核心1 k, err := createAndInitKubelet(kubeServer, kubeDeps, hostname, hostnameOverridden, nodeName, nodeIPs) ... // 核心2 startKubelet(k, podCfg, &kubeServer.KubeletConfiguration, kubeDeps, kubeServer.EnableServer) ... }1
2
3
4
5
6
7
8
9
10
11
12
13
14createAndInitKubelet()初始化kubelet实例,启动GC回收器件,同时向apiserver发送kubelet start event。func createAndInitKubelet(...) (k kubelet.Bootstrap, err error) { // 初始化kubelet实例 k, err = kubelet.NewMainKubelet(...) ... // 事件上报 k.BirthCry() // 启动GC k.StartGarbageCollection() return k, nil } // NewMainKubelet instantiates a new Kubelet object along with all the required internal modules. func NewMainKubelet(...) (*Kubelet, error) { ... // 节点信息同步(List Watch) var nodeHasSynced cache.InformerSynced var nodeLister corelisters.NodeLister ... nodeLister = kubeInformers.Core().V1().Nodes().Lister() nodeHasSynced = func() bool { return kubeInformers.Core().V1().Nodes().Informer().HasSynced() } kubeInformers.Start(wait.NeverStop) ... // service同步 var serviceLister corelisters.ServiceLister var serviceHasSynced cache.InformerSynced ... serviceLister = kubeInformers.Core().V1().Services().Lister() serviceHasSynced = kubeInformers.Core().V1().Services().Informer().HasSynced kubeInformers.Start(wait.NeverStop) ... // 初始化oomWatcher oomWatcher, err := oomwatcher.NewWatcher(kubeDeps.Recorder) ... // 解析集群DNS配置 clusterDNS := make([]net.IP, 0, len(kubeCfg.ClusterDNS)) for _, ipEntry := range kubeCfg.ClusterDNS { ip := netutils.ParseIPSloppy(ipEntry) clusterDNS = append(clusterDNS, ip) } ... // 初始化kubelet klet := &Kubelet{} ... // 配置Secret和ConfigMap管理器(watch/ttl/get) ... klet.secretManager = secretManager klet.configMapManager = configMapManager ... // 存活探针管理器 klet.livenessManager = proberesults.NewManager() // 就绪探针管理器 klet.readinessManager = proberesults.NewManager() // 启动探针管理器 klet.startupManager = proberesults.NewManager() // Pod缓存管理器 klet.podCache = kubecontainer.NewCache() ... // Pod管理器 klet.podManager = kubepod.NewBasicPodManager(mirrorPodClient, secretManager, configMapManager) // 状态管理器 klet.statusManager = status.NewManager(klet.kubeClient, klet.podManager, klet) ... // 运行时管理器 runtime, err := kuberuntime.NewKubeGenericRuntimeManager( kubecontainer.FilterEventRecorder(kubeDeps.Recorder), klet.livenessManager, // 存活探针管理器 klet.readinessManager, // 就绪探针管理器 klet.startupManager, // 启动探针管理器 ... klet.podWorkers, // pod变更管理器 ... kubeDeps.RemoteRuntimeService, // 运行时服务 kubeDeps.RemoteImageService, // 镜像服务 kubeDeps.ContainerManager.InternalContainerLifecycle(), // 容器生命周期管理器 klet.containerLogManager, // 容器日志管理器 ... ) ... klet.containerRuntime = runtime klet.streamingRuntime = runtime ... // 运行时缓存 klet.runtimeCache = runtimeCache ... // 统计管理(基于cAdvisor/CRI) klet.StatsProvider = stats.NewCadvisorStatsProvider()/stats.NewCRIStatsProvider() ... // pod生命周期事件生成器 klet.pleg = pleg.NewGenericPLEG(...) ... // 容器GC klet.containerGC = containerGC klet.containerDeletor = newPodContainerDeletor(...) ... // 镜像GC klet.imageManager = imageManager // 证书管理器 klet.serverCertificateManager, err = kubeletcertificate.NewKubeletServerCertificateManager(...) ... // 探测管理器 klet.probeManager = prober.NewManager(...) ... // 卷插件管理器 klet.volumePluginMgr = NewInitializedVolumePluginMgr() ... // 插件管理器(GPU/FPGA) klet.pluginManager = pluginmanager.NewPluginManager(...) ... // 卷管理器 klet.volumeManager = volumemanager.NewVolumeManager(...) ... // 驱逐管理器 klet.evictionManager = evictionManager klet.admitHandlers.AddPodAdmitHandler(evictionAdmitHandler) ... // 节点租约控制器(ready/notReady) klet.nodeLeaseController = lease.NewController(...) ... // 节点管理管理器 klet.shutdownManager = shutdownManager klet.admitHandlers.AddPodAdmitHandler(shutdownAdmitHandler) // Finally, put the most recent version of the config on the Kubelet, so // people can see how it was configured. klet.kubeletConfiguration = *kubeCfg // Node默认信息注入(address/machineInfo) klet.setNodeStatusFuncs = klet.defaultNodeStatusFuncs() return klet, nil } // GC回收期 func (kl *Kubelet) StartGarbageCollection() { go wait.Until(func() { // 容器GC回收 kl.containerGC.GarbageCollect() ... }, ContainerGCPeriod, wait.NeverStop) go wait.Until(func() { // 镜像GC回收 kl.imageManager.GarbageCollect() ... }, ImageGCPeriod, wait.NeverStop) }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
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146startKubelet以协程启动kubelet中的各个模块,注册kubelet http server。--- 启一个协程调用Run()方法(kubelet核心逻辑) func startKubelet(...) { // 启动kubelet go k.Run(podCfg.Updates()) // kubectl http server go k.ListenAndServe(kubeCfg, kubeDeps.TLSOptions, kubeDeps.Auth) ... }1
2
3
4
5
6
7
8
9
# 3.3.kubelet.Run
Run方法启动kubelet的依赖模块,尤其是主循环逻辑,用于各模块启动和pod生命周期管理、外部信息同步。// Run starts the kubelet reacting to config updates func (kl *Kubelet) Run(updates <-chan kubetypes.PodUpdate) { // 注册LogServer kl.logServer = http.StripPrefix("/logs/", http.FileServer(http.Dir("/var/log/"))) ... // 初始化不依赖container runtime的模块 kl.initializeModules() // 启动volume manager,处理挂载/卸载卷 go kl.volumeManager.Run(kl.sourcesReady, wait.NeverStop) // 定期同步Node状态 go wait.JitterUntil(kl.syncNodeStatus, kl.nodeStatusUpdateFrequency, 0.04, true, wait.NeverStop) // 更新容器运行时启动时间以及执行首次状态同步 go kl.fastStatusUpdateOnce() // 同步当前node的kubelet租约(心跳) go kl.nodeLeaseController.Run(wait.NeverStop) // 每个5s检查一次容器运行时状态 go wait.Until(kl.updateRuntimeUp, 5*time.Second, wait.NeverStop) // 初始化iptables配置链(svc访问apiserver) if kl.makeIPTablesUtilChains { kl.initNetworkUtil() } // 启动状态管理器,同步pod状态到apiserver kl.statusManager.Start() // 启动运行时类同步(管理不同运行时) if kl.runtimeClassManager != nil { kl.runtimeClassManager.Start(wait.NeverStop) } // 定期检查容器运行时,生成pod生命周期事件 kl.pleg.Start() // 处理pod变更事件 kl.syncLoop(updates, kl) }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
# 3.4.initializeModules
// initializeModules will initialize internal modules that do not require the container runtime to be up. func (kl *Kubelet) initializeModules() error { ... // 初始化/var/lib/kubelet、/var/lib/kubelet/pods、/var/lib/kubelet/plugins目录 kl.setupDataDirs() // 初始化/var/log/containers kl.os.MkdirAll(ContainerLogsDir, 0755) ... // 启动镜像管理器(维护未使用镜像记录,用于GC) kl.imageManager.Start() // 启动证书管理器(证书轮转) kl.serverCertificateManager.Start() ... // 启动oomWatcher(监听/dev/kmsg的oom kill事件) kl.oomWatcher.Start(kl.nodeRef) ... return nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18kl.oomWatcher.Start启动oomWatcher监听器,主要监听/dev/kmsg文件日志获取未关联container的oom事件记录为event。// Start watches for system oom's and records an event for every system oom encountered. func (ow *realWatcher) Start(ref *v1.ObjectReference) error { outStream := make(chan *oomparser.OomInstance, 10) // 监听节点上OOM事件 go ow.oomStreamer.StreamOoms(outStream) go func() { ... // 遍历每个事件 for event := range outStream { // 系统级别的OOM(未关联容器的OOM事件) if event.VictimContainerName == recordEventContainerName { ... // 记录OOM事件至APIServer(Node级别) ow.recorder.Eventf(ref, v1.EventTypeWarning, systemOOMEvent, eventMsg) } } }() return nil } // StreamOoms writes to a provided a stream of OomInstance objects representing // OOM events that are found in the logs. // It will block and should be called from a goroutine. func (p *OomParser) StreamOoms(outStream chan<- *OomInstance) { // 循环读取/dev/kmsg的日志到kmsgEntries管道 kmsgEntries := p.parser.Parse() ... for msg := range kmsgEntries { // 检查是否属于OOM事件 isOomMessage := checkIfStartOfOomMessages(msg.Message) if isOomMessage { // 初始化OOM实例 oomCurrentInstance := &OomInstance{ ContainerName: "/", VictimContainerName: "/", TimeOfDeath: msg.Timestamp, } // 解析多行日志的容器名及PID for msg := range kmsgEntries { // 获取容器名 finished, err := getContainerName(msg.Message, oomCurrentInstance) ... // 找不到容器名获取PID finished, err = getProcessNamePid(msg.Message, oomCurrentInstance) if finished { oomCurrentInstance.TimeOfDeath = msg.Timestamp break } } // 推送到事件管道 outStream <- oomCurrentInstance } } } // Parse will read from the provided reader and provide a channel of messages func (p *parser) Parse() <-chan Message { // OOM事件管道 output := make(chan Message, 1) go func() { defer close(output) msg := make([]byte, 8192) for { // 读取/dev/kmsg日志 n, err := p.kmsgReader.Read(msg) ... msgStr := string(msg[:n]) // 解析日志为OOM消息 message, err := p.parseMessage(msgStr) ... output <- message } }() return output }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
# 4.syncLoop
# 4.1.syncLoop
syncLoop是kubelet的主循环方法,根据不同数据来源——file、url及apiserver监听pod变化,调度对应函数处理,保证pod处于期望状态。func (kl *Kubelet) syncLoop(updates <-chan kubetypes.PodUpdate, handler SyncHandler) { ... for { ... // 反复调用syncLoopIterration处理Pod变更 if !kl.syncLoopIteration(updates, handler, syncTicker.C, housekeepingTicker.C, plegCh) { break } } }1
2
3
4
5
6
7
8
9
10
# 4.2.syncLoopIteration
syncLoopIteration监听多个管道事件,调用相应的handler-->dispatchWork分发任务。func (kl *Kubelet) syncLoopIteration(...) bool { select { // 1.watch不同来源的pod信息变化(file、http、apiserver) case u, open := <-configCh: ... switch u.Op { case kubetypes.ADD: // 新增 handler.HandlePodAdditions(u.Pods) case kubetypes.UPDATE: // 更新 handler.HandlePodUpdates(u.Pods) case kubetypes.REMOVE: // 移除 handler.HandlePodRemoves(u.Pods) case kubetypes.RECONCILE: // 重新协调 handler.HandlePodReconcile(u.Pods) case kubetypes.DELETE: // 优雅删除 handler.HandlePodUpdates(u.Pods) ... } // 2.pleg.start()每秒reList容器状态,根据最新的PodStatus生成PodLifeCycleEvent存入PLEChan case e := <-plegCh: // 更新容器最后一次启动时间 if e.Type == pleg.ContainerStarted { kl.lastContainerStartedTime.Add(e.ID, time.Now()) } // 容器状态更新 if isSyncPodWorthy(e) { handler.HandlePodSyncs([]*v1.Pod{pod}) } // 容器退出 if e.Type == pleg.ContainerDied { if containerID, ok := e.Data.(string); ok { kl.cleanUpContainersInPod(e.ID, containerID) } } // 3.每秒周期执行 case <-syncCh: // 获取所有待同步的pod(运行正常的pod/内部模块请求的pod) podsToSync := kl.getPodsToSync() ... // 同步最新的Pod状态 handler.HandlePodSyncs(podsToSync) // 4.liveness事件处理 case update := <-kl.livenessManager.Updates(): // 如果探针检测失败,触发探测同步 if update.Result == proberesults.Failure { handleProbeSync(kl, update, handler, "liveness", "unhealthy") } // 5.readiness事件处理 case update := <-kl.readinessManager.Updates(): // 获取readiness探测状态 ready := update.Result == proberesults.Success // 更新pod状态 kl.statusManager.SetContainerReadiness(update.PodUID, update.ContainerID, ready) if ready { status = "ready" } // 触发探测同步 handleProbeSync(kl, update, handler, "readiness", status) // 6.启动状态变化 case update := <-kl.startupManager.Updates(): started := update.Result == proberesults.Success // 更新容器状态 kl.statusManager.SetContainerStartup(update.PodUID, update.ContainerID, started) if started { status = "started" } // 触发探测同步 handleProbeSync(kl, update, handler, "startup", status) // 7.每2秒钟执行一次 case <-housekeepingCh: // 所有配置源就绪才触发 if kl.sourcesReady.AllReady() { // 家务进程,清理孤儿容器、回收卷、清理终止pod handler.HandlePodCleanups() } } return true }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
事件来源
1.
configCh接收不同来源pod变化,同步file/http/apiserver的pod变化事件2.
plegCh接收pleg模块同步的pod状态后生成的事件,通过handler中进行dispatchWork分发reList的相关事件3.
syncCh接收所有等待同步的pod,每秒进行处理4.
livenessCh接收所有的liveness探测事件,对失败的或者liveness检查失败的pod进行同步5.
readinessCh接收所有的readiness探测事件,对失败的或readiness检查失败的pod进行同步6.
housekeepingCh接收watch的清理事件,每2秒钟触发对pod的清理