hpa
# 1.简介
# 1.1.定义
hpa意为自动伸缩,基于*metrics.k8s.io APISerice访问各种metrics提供者,基于CPU/Mem等动态伸缩Pod的副本数量。// HorizontalController is responsible for the synchronizing HPA objects stored in the system. type HorizontalController struct { scaleNamespacer scaleclient.ScalesGetter // scale client hpaNamespacer autoscalingclient.HorizontalPodAutoscalersGetter // hpa client mapper apimeta.RESTMapper // kind-->GV导航 replicaCalc *ReplicaCalculator // replica计算器 ... downscaleStabilisationWindow time.Duration // 典型值5min,抑制缩容过快 monitor monitor.Monitor // 底层metric采集的抽象接口 hpaLister autoscalinglisters.HorizontalPodAutoscalerLister // hpa informer缓存 ... podLister corelisters.PodLister // Pod informer缓存 ... queue workqueue.RateLimitingInterface // hpa工作队列 recommendations map[string][]timestampedRecommendation // hpa的历史推荐副本数 ... scaleUpEvents map[string][]timestampedScaleEvent // hpa的历史扩容行为 ... scaleDownEvents map[string][]timestampedScaleEvent // hpa的历史缩容行为 ... // Storage of HPAs and their selectors. hpaSelectors *selectors.BiMultimap // hpa<-->Pod双向selector映射 ... }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
注意
hpa本身不直接管理Pod,而是管理deployment/replicaset的副本数实现Pod动态伸缩
# 1.2.原理
hpa根据目标的资源使用率计算Pod合理的扩缩上限,通常修改deployment/replicaset的副本数完成Pod数量动态调整,实现负载调整。
# 2.分析
# 2.1.start
startHPAController()会检查hpa资源启用状态,调用startHPAControllerWithRESTClient()实例化及启动hpa controller。func startHPAController(...) (controller.Interface, bool, error) { // 未启用HPA if !controllerContext.AvailableResources[schema.GroupVersionResource{Group: "autoscaling", Version: "v1", Resource: "horizontalpodautoscalers"}] { return nil, false, nil } return startHPAControllerWithRESTClient(ctx, controllerContext) } func startHPAControllerWithRESTClient(...) (controller.Interface, bool, error) { ... // 初始化apiVersionsFromDiscovery apiVersionsGetter := custom_metrics.NewAvailableAPIsGetter(hpaClient.Discovery()) // 间隔15s重置apiVersionsFromDiscovery的perfVersion缓存 go custom_metrics.PeriodicallyInvalidate(apiVersionsGetter, 15s, ctx.Done()) // 初始化Pod metric client metricsClient := metrics.NewRESTMetricsClient( // 内置metrics查询(apiservice对应metrics server) resourceclient.NewForConfigOrDie(clientConfig), // custom metrics查询 custom_metrics.NewForConfig(clientConfig, controllerContext.RESTMapper, apiVersionsGetter), // external metrics查询 external_metrics.NewForConfigOrDie(clientConfig), ) // 实例化及启动hpa return startHPAControllerWithMetricsClient(ctx, controllerContext, metricsClient) } func startHPAControllerWithMetricsClient(...) (controller.Interface, bool, error) { ... // 实例化hpa及激活 go podautoscaler.NewHorizontalController( hpaClient.CoreV1(), // event client scaleClient, // scale client hpaClient.AutoscalingV2(), // hpa client controllerContext.RESTMapper, // kind<-->GV映射 metricsClient, // metric client controllerContext.InformerFactory.Autoscaling().V2().HorizontalPodAutoscalers(), // hpa informer缓存 controllerContext.InformerFactory.Core().V1().Pods(), // Pod informer缓存 ... ).Run(ctx, int(controllerContext.ComponentConfig.HPAController.ConcurrentHorizontalPodAutoscalerSyncs)) return nil, true, 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
37
38
39
40
41
42
43
44
45
46
47
48补充
metricsClient作为metrics的聚合客户端,获取及缓存resource/custom/external metrics
# 2.2.informer
hpa controller基于informer监听及回调hpa资源、缓存Pod及构造副本计算器replicaCalc,replicaCalc根据指标计算预期副本数。// NewHorizontalController creates a new HorizontalController. func NewHorizontalController(...) *HorizontalController { ... hpaController := &HorizontalController{ ... scaleNamespacer: scaleNamespacer, hpaNamespacer: hpaNamespacer, downscaleStabilisationWindow: downscaleStabilisationWindow, // 缩容间隔5min monitor: monitor.New(), // 延迟15s入队 queue: workqueue.NewNamedRateLimitingQueue(NewDefaultHPARateLimiter(resyncPeriod), "hpa"), mapper: mapper, recommendations: map[string][]timestampedRecommendation{}, .... scaleUpEvents: map[string][]timestampedScaleEvent{}, ... scaleDownEvents: map[string][]timestampedScaleEvent{}, ... hpaSelectors: selectors.NewBiMultimap(), ... } // hpa informer(间隔15s同步一次) hpaInformer.Informer().AddEventHandlerWithResyncPeriod( cache.ResourceEventHandlerFuncs{ AddFunc: hpaController.enqueueHPA, UpdateFunc: hpaController.updateHPA, DeleteFunc: hpaController.deleteHPA, }, resyncPeriod, ) hpaController.hpaLister = hpaInformer.Lister() ... // Pod informer(仅Lister查询) hpaController.podLister = podInformer.Lister() ... // 副本计算器 replicaCalc := NewReplicaCalculator( metricsClient, hpaController.podLister, tolerance, // 副本差异容忍范围(0.1) cpuInitializationPeriod, // CPU初始化周期(5min) delayOfInitialReadinessStatus, // 就绪延迟(30s) ) hpaController.replicaCalc = replicaCalc ... return hpaController } // obj could be an *v1.HorizontalPodAutoscaler, or a DeletionFinalStateUnknown marker item. func (a *HorizontalController) enqueueHPA(obj interface{}) { key, err := controller.KeyFunc(obj) ... // 延迟15s入队 a.queue.AddRateLimited(key) ... // hpa的label selector先初始化为空 if hpaKey := selectors.Parse(key); !a.hpaSelectors.SelectorExists(hpaKey) { a.hpaSelectors.PutSelector(hpaKey, labels.Nothing()) } } // obj could be an *v1.HorizontalPodAutoscaler, or a DeletionFinalStateUnknown marker item. func (a *HorizontalController) updateHPA(old, cur interface{}) { a.enqueueHPA(cur) } func (a *HorizontalController) deleteHPA(obj interface{}) { key, err := controller.KeyFunc(obj) ... // 重置重试状态 a.queue.Forget(key) ... // 清理label selector a.hpaSelectors.DeleteSelector(selectors.Parse(key)) }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注意
hpaInformer会监听资源及执行入队回调,podInformer仅缓存pod数据
# 2.3.runworker
hpa.Run()会激活worker消费workqueue的任务,workqueue基于固定时间参数(15s)限速,实现类似ticker的定时执行效果。// Run begins watching and syncing. func (a *HorizontalController) Run(ctx context.Context, workers int) { ... defer a.queue.ShutDown() ... // informer同步完成 if !cache.WaitForNamedCacheSync("HPA", ctx.Done(), a.hpaListerSynced, a.podListerSynced) { return } // 激活5个worker for i := 0; i < workers; i++ { go wait.UntilWithContext(ctx, a.worker, time.Second) } <-ctx.Done() } func (a *HorizontalController) worker(ctx context.Context) { for a.processNextWorkItem(ctx) { } ... } func (a *HorizontalController) processNextWorkItem(ctx context.Context) bool { ... defer a.queue.Done(key) // 动态伸缩 deleted, err := a.reconcileKey(ctx, key.(string)) ... // 无异常且未被删除,重新延迟15s入队 if !deleted { a.queue.AddRateLimited(key) } 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注意
hpa未删除及处理完成,会延迟15s延迟入队(基于heap),等待下一轮时间周期会再次取出执行
# 3.入口
# 3.1.reconcile
hpa.reconcileKey()会获取hpa对象,hpa不存在会清理相关的缓存数据,否则执行hpa.reconcileAutoscaler()进行副本评估及扩缩容。func (a *HorizontalController) reconcileKey(ctx context.Context, key string) (deleted bool, err error) { namespace, name, err := cache.SplitMetaNamespaceKey(key) ... // 获取hpa对象 hpa, err := a.hpaLister.HorizontalPodAutoscalers(namespace).Get(name) // hpa对象未找到 if k8serrors.IsNotFound(err) { ... // hpa的历史副本数 delete(a.recommendations, key) ... // hpa历史扩容行为 delete(a.scaleUpEvents, key) ... // hpa历史缩容行为 delete(a.scaleDownEvents, key) ... return true, nil } ... // 执行scale流程 return false, a.reconcileAutoscaler(ctx, hpa, key) }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注意
replica计算及执行scale的核心逻辑都由a.reconcileAutoscaler()完成
# 3.2.autoscaler
a.reconcileAutoscaler()根据当前副本数和HPA对象定义的期望副本数计算调整的数量,基于scale执行扩容或缩容,这是HPA的核心方法。func (a *HorizontalController) reconcileAutoscaler(...) (retErr error) { ... hpa := hpaShared.DeepCopy() hpaStatusOriginal := hpa.Status.DeepCopy() ... // ref的GV targetGV, err := schema.ParseGroupVersion(hpa.Spec.ScaleTargetRef.APIVersion) ... // ref的GK targetGK := schema.GroupKind{ Group: targetGV.Group, Kind: hpa.Spec.ScaleTargetRef.Kind, } // 基于GK获取rest信息 mappings, err := a.mapper.RESTMappings(targetGK) ... // 基于rest读取scale子资源 scale, targetGR, err := a.scaleForResourceMappings(ctx,hpa.Namespace,hpa.Spec.ScaleTargetRef.Name,mappings) ... // 更新condition setCondition(hpa, autoscalingv2.AbleToScale, v1.ConditionTrue, "SucceededGetScale", ...) // 获取当前副本数 currentReplicas := scale.Spec.Replicas // 缓存hpa对象初始化副本 a.recordInitialRecommendation(currentReplicas, key) ... // 设置最小副本数 if hpa.Spec.MinReplicas != nil { minReplicas = *hpa.Spec.MinReplicas } else { // Default value minReplicas = 1 } ... // curReplicas为0但minReplicas不为0 if scale.Spec.Replicas == 0 && minReplicas != 0 { // 禁用自动扩缩容 desiredReplicas = 0 rescale = false setCondition(hpa, autoscalingv2.ScalingActive, v1.ConditionFalse, "ScalingDisabled", ...) // curReplicas超过maxReplicas } else if currentReplicas > hpa.Spec.MaxReplicas { // desiredReplicas设为maxReplicas rescaleReason = "Current number of replicas above Spec.MaxReplicas" desiredReplicas = hpa.Spec.MaxReplicas // desiredReplicas低于minReplicas } else if currentReplicas < minReplicas { // 重置为minReplicas desiredReplicas = minReplicas } else { ... // 基于metrics计算期望副本数 mdr, metricName, mStatus, mts, err = a.computeReplicasForMetrics(ctx, hpa, scale, hpa.Spec.Metrics) ... // metrics计算出的replicas大一些 if mdr > desiredReplicas { // 更新desiredReplicas desiredReplicas = mdr // 记录相关的指标名称 rescaleMetric = metricName } ... // 基于扩缩窗口调整 if hpa.Spec.Behavior == nil { desiredReplicas = a.normalizeDesiredReplicas(hpa,key,currentReplicas,desiredReplicas,minReplicas) // 基于behavior调整 } else { desiredReplicas = a.normalizeDesiredReplicasWithBehaviors(hpa,key,currentReplicas,desiredReplicas, minReplicas) } // 检查是否scale rescale = desiredReplicas != currentReplicas } // 触发scale if rescale { // 预期副本数 scale.Spec.Replicas = desiredReplicas // 执行scale更新 _, err = a.scaleNamespacer.Scales(hpa.Namespace).Update(ctx, targetGR, scale, metav1.UpdateOptions{}) ... // 缓存scale行为 a.storeScaleEvent(hpa.Spec.Behavior, key, currentReplicas, desiredReplicas) ... // 未触发scale } else { // desiredReplicas重置为当前副本数 desiredReplicas = currentReplicas } // 更新status a.setStatus(hpa, currentReplicas, desiredReplicas, mStatus, rescale) a.updateStatusIfNeeded(ctx, hpaStatusOriginal, hpa) ... return retErr }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注意
curReplicas位于minReplicas和maxReplicas之间会基于metrics计算推荐副本,否则会尝试设置为minReplicas或maxReplicas
# 3.3.computeR
a.computeReplicasForMetrics()根据HPA对象配置的metric计算预期的副本数,计算时取最大值作为预期副本,以覆盖所有资源压力。// computeReplicasForMetrics computes the desired replicas for the metric specifications listed in the HPA. func (a *HorizontalController) computeReplicasForMetrics(...) (...) { // scale selector检查(基于selector可匹配hpa和Pod) selector, err := a.validateAndParseSelector(hpa, scale.Status.Selector) ... // 遍历配置的指标 for i, metricSpec := range metricSpecs { // 基于metric计算预期副本 rc, metricName, ts, condition, err := a.computeReplicasForMetric(ctx, hpa, metricSpec, specReplicas, statusReplicas, selector, &statuses[i]) ... // 更新为更大的副本数 if replicas == 0 || rc > replicas { timestamp = ts replicas = rc metric = metricName } } ... // 设置condition setCondition(hpa, autoscalingv2.ScalingActive, v1.ConditionTrue, "ValidMetricFound", ...) return replicas, metric, statuses, timestamp, invalidMetricError } // Computes the desired replicas for a specific hpa and metric specification. func (a *HorizontalController) computeReplicasForMetric(...) (...) { ... switch spec.Type { // 基于k8s内置对象的指标计算 case autoscalingv2.ObjectMetricSourceType: ... rc, ts, metricName, condition, err = a.computeStatusForObjectMetric(specReplicas, statusReplicas, spec, hpa, selector, status, metricSelector) ... // 基于Pod metric计算 case autoscalingv2.PodsMetricSourceType: ... rc, ts, metricName, condition, err = a.computeStatusForPodsMetric(specReplicas, spec, hpa, selector, status, metricSelector) ... // 基于CPU/Mem资源指标计算 case autoscalingv2.ResourceMetricSourceType: rc, ts, metricName, condition, err = a.computeStatusForResourceMetric(ctx, specReplicas, spec, hpa, selector, status) ... // 基于container的资源指标计算 case autoscalingv2.ContainerResourceMetricSourceType: ... rc, ts, metricName, condition, err = a.computeStatusForContainerResourceMetric(ctx, specReplicas, spec, hpa, selector, status) ... // 基于external metric计算 case autoscalingv2.ExternalMetricSourceType: replicaCountProposal, timestampProposal, metricNameProposal, condition, err = a.computeStatusForExternalMetric(specReplicas, statusReplicas, spec, hpa, selector, status) ... } return rc, metricName, ts, autoscalingv2.HorizontalPodAutoscalerCondition{}, 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
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60注意
hpa基于不同metric计算期望副本执行scale,常用的是resourceMetric