nodeipam
# 1.简介
# 1.1.定义
nodeipam controller用于向node节点分配不冲突的可用cidr,kubelet基于node关联的cidr为Pod申请addr。
注意
kubelet基于CRI创建Pod会交互CNI网络插件,基于node cidr分配Pod地址,flannel那里分析过
# 1.2.原理
nodeipam controller基于cidr信息生成网段分配器,监听node对象变化,ADD事件会向node分配cidr,DEL事件会释放node cidr。
注意
nodeipam仅负责划分node cidr供CNI分配,service clusterIP cidr分配由apiserver负责
# 2.入口
# 2.1.start
startNodeIpamController()会检查cluster cidr相关前置条件,实例化nodeIpamController及启动,监听node及触发相关事件回调。func startNodeIpamController(...) (controller.Interface, bool, error) { ... // 解析clusterCIDR列表(单栈1个,双栈2个) clusterCIDRs, err := validateCIDRs(controllerContext.ComponentConfig.KubeCloudShared.ClusterCIDR) ... // 解析service cidr if len(strings.TrimSpace(controllerContext.ComponentConfig.NodeIPAMController.ServiceCIDR)) != 0 { _, serviceCIDR, err = netutils.ParseCIDRSloppy(...NodeIPAMController.ServiceCIDR) ... } if len(strings.TrimSpace(controllerContext.ComponentConfig.NodeIPAMController.SecondaryServiceCIDR)) != 0 { _, secondaryServiceCIDR, err = netutils.ParseCIDRSloppy(...NodeIPAMController.SecondaryServiceCIDR) ... } // 单栈1个service cidr,双栈2个 if serviceCIDR != nil && secondaryServiceCIDR != nil { // should be dual stack (from different IPFamilies) dualstackServiceCIDR, err := netutils.IsDualStackCIDRs([]*net.IPNet{serviceCIDR, secondaryServiceCIDR}) ... if !dualstackServiceCIDR { return nil, false, err } } // maskSize解析及clusterCIDR分组 nodeCIDRMaskSizes, err := setNodeCIDRMaskSizes(...NodeIPAMController, clusterCIDRs) ... // clusterCIDR informer监听 var clusterCIDRInformer v1alpha1.ClusterCIDRInformer if utilfeature.DefaultFeatureGate.Enabled(features.MultiCIDRRangeAllocator) { clusterCIDRInformer = controllerContext.InformerFactory.Networking().V1alpha1().ClusterCIDRs() } // 实例化nipc nodeIpamController, err := nodeipamcontroller.NewNodeIpamController( ctx, controllerContext.InformerFactory.Core().V1().Nodes(), clusterCIDRInformer, ... clusterCIDRs, serviceCIDR, secondaryServiceCIDR, nodeCIDRMaskSizes, ipam.CIDRAllocatorType(controllerContext.ComponentConfig.KubeCloudShared.CIDRAllocatorType), ) ... // 激活 go nodeIpamController.RunWithMetrics(ctx, controllerContext.ControllerManagerMetrics) 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
49
50
51
52
53
54
55
56
57
注意
startNodeIpamController()会检查clusterCIDR及serviceCIDR有效性,基于cidr实例化nipc及启动
# 2.2.new
NewNodeIpamController()会检查cidr/mask合法性,基于不同的分配模式初始化不同的allocator,激活nodeInformer及返回nipc。func createLegacyIPAM(...) (*ipam.Controller, error) { cfg := &ipam.Config{ Resync: ipamResyncInterval, // 间隔30s同步一次 MaxBackoff: ipamMaxBackoff, // 最大重试10s InitialRetry: ipamInitialBackoff, // 初始重试250ms } // 设置分同步模式 switch ic.allocatorType { case ipam.IPAMFromClusterAllocatorType: cfg.Mode = nodesync.SyncFromCluster case ipam.IPAMFromCloudAllocatorType: cfg.Mode = nodesync.SyncFromCloud } // we may end up here with no cidr at all in case of FromCloud/FromCluster var cidr *net.IPNet if len(clusterCIDRs) > 0 { cidr = clusterCIDRs[0] } ... // 分配器 ipamc, err := ipam.NewController(cfg, kubeClient, cloud, cidr, serviceCIDR, nodeCIDRMaskSizes[0]) ... // 定期同步 ipamc.Start(logger, nodeInformer) ... return ipamc, nil } // New creates a new CIDR range allocator. func New(...) (CIDRAllocator, error) { // 获取node列表 nodeList, err := listNodes(logger, kubeClient) ... switch allocatorType { // Range类型分配器 case RangeAllocatorType: return NewCIDRRangeAllocator(logger, kubeClient, nodeInformer, allocatorParams, nodeList) // MultiCIDR分配器 case MultiCIDRRangeAllocatorType: ... return NewMultiCIDRRangeAllocator(ctx, kubeClient, nodeInformer, clusterCIDRInformer, allocatorParams, nodeList, nil) // Cloud类型分配器 case CloudAllocatorType: return NewCloudCIDRAllocator(logger, kubeClient, cloud, nodeInformer) ... } } // returns a new node IP Address Management controller to sync instances from cloudprovider. func NewNodeIpamController(...) (*Controller, error) { ... // 非cloud cidrAllocator if allocatorType != ipam.CloudAllocatorType { // 检查cidr if len(clusterCIDRs) == 0 { return nil, fmt.Errorf("Controller: Must specify --cluster-cidr if --allocate-node-cidrs is set") } // 检查maskSize for idx, cidr := range clusterCIDRs { mask := cidr.Mask if maskSize, _ := mask.Size(); maskSize > nodeCIDRMaskSizes[idx] { return nil, err } } } // 实例化ipam controller ic := &Controller{ cloud: cloud, kubeClient: kubeClient, eventBroadcaster: record.NewBroadcaster(), lookupIP: net.LookupIP, clusterCIDRs: clusterCIDRs, serviceCIDR: serviceCIDR, secondaryServiceCIDR: secondaryServiceCIDR, allocatorType: allocatorType, } // 仅同步不参与分配 if ic.allocatorType == IPAMFromCluster || ic.allocatorType == IPAMFromCloud { // 同步 ic.legacyIPAM, _ = createLegacyIPAM(logger, ic, nodeInformer, cloud, kubeClient, clusterCIDRs, serviceCIDR, nodeCIDRMaskSizes) ... } else { ... // 构造分配器 ic.cidrAllocator, err = ipam.New(ctx, kubeClient, cloud, nodeInformer, clusterCIDRInformer, ic.allocatorType, allocatorParams) ... } ic.nodeLister = nodeInformer.Lister() ic.nodeInformerSynced = nodeInformer.Informer().HasSynced return ic, 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
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
注意
nipc实例化会基于allocatorType初始化cidr分配器,变革类型涉及range/multi/cloud多种,后续仅分析rangeAllocator
# 2.3.run
nc.RunWithMetrics()会激活nodeIpamController初始化的分配器,触发legacyIPAM或cidrAllocator启动执行cidr分配。// wrapper for Run that also tracks starting and stopping of the nodeipam controller with additional metric func (nc *Controller) RunWithMetrics(...) { ... nc.Run(ctx) } // Run starts an asynchronous loop that monitors the status of cluster nodes. func (nc *Controller) Run(ctx context.Context) { ... // node同步等待 if !cache.WaitForNamedCacheSync("node", ctx.Done(), nc.nodeInformerSynced) { return } // 传统分配器 if nc.allocatorType == IPAMFromCluster || nc.allocatorType == IPAMFromCloud { go nc.legacyIPAM.Run(ctx) // 新类型分配器 } else { go nc.cidrAllocator.Run(ctx) } <-ctx.Done() } // legacyIPAM分配器,什么都不做,仅阻塞 func (c *Controller) Run(ctx context.Context) { ... go c.adapter.Run(ctx) <-ctx.Done() } func (a *adapter) Run(ctx context.Context) { ... <-ctx.Done() }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
注意
nc.legacyIPAM.Run()本质什么都不做,分配或同步工作已基于syncer完成,nc.cidrAllocator.Run()激活worker执行不同分配策略
# 3.legacy
# 3.1.new
createLegacyIPAM()会执行ipam.NewController()实例化分配器,调用ipamc.Start()周期同步cluster/cloud的cidr至node。// NewController returns a new instance of the IPAM controller. func NewController(...) (*Controller, error) { ... // 仅适配GCE云服务商 gceCloud, ok := cloud.(*gce.Cloud) ... // 基于mask构造cidrSet // nodeMask和cidrMask最大差16,避免造成地址浪费 set, err := cidrset.NewCIDRSet(clusterCIDR, nodeCIDRMaskSize) ... // controller实例化 c := &Controller{ config: config, adapter: newAdapter(kubeClient, gceCloud), syncers: make(map[string]*nodesync.NodeSync), set: set, } // cidr规避serviceCIDR重合部分 occupyServiceCIDR(c.set, clusterCIDR, serviceCIDR) ... // 向后申请一个cidr(检查剩余) cidr, err := c.set.AllocateNext() switch err { ... case nil: // 释放cidr(begin~end均释放) err := c.set.Release(cidr) return c, err ... } }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
注意
legacy ipam是传统分配器,基于cluster-cidr或cloud alias分配或同步node cidr
# 3.2.start
ipamc.Start()会注册node资源变更回调,周期同步cluster/cloud的cidr更新至node,这也是传统分配模式确保数据一致性的手段。// initializes the Controller with the existing list of nodes and registers the informers for node changes. func (c *Controller) Start(logger klog.Logger, nodeInformer informers.NodeInformer) error { // 获取node列表 nodes, err := listNodes(logger, c.adapter.k8s) ... // 遍历检查 for _, node := range nodes.Items { // node cidr不为空 if node.Spec.PodCIDR != "" { // 解析cidr _, cidrRange, err := netutils.ParseCIDRSloppy(node.Spec.PodCIDR) ... // 标记此cidr地址已占用 c.set.Occupy(cidrRange) } func() { ... // 注册nodeSync及启动 syncer := c.newSyncer(node.Name) c.syncers[node.Name] = syncer go syncer.Loop(logger, nil) }() } // node informer监听 nodeInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{ AddFunc: controllerutil.CreateAddNodeHandler(func(node *v1.Node) error { return c.onAdd(logger, node) }), UpdateFunc: controllerutil.CreateUpdateNodeHandler(func(_, newNode *v1.Node) error { return c.onUpdate(logger, newNode) }), DeleteFunc: controllerutil.CreateDeleteNodeHandler(func(node *v1.Node) error { return c.onDelete(logger, node) }), }) 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
37
38
39
40
41
42
43
补充
nodeInformer仅作为事件来源,驱动sync及时同步node cidr
# 3.3.syncer
nodesyncer注册完成会执行sync.Loop()启动异步循环,间隔30s触发一次sync.opChan消费,获取及执行node cidr更新或清理。// Loop runs the sync loop for a given node. done is an optional channel that is closed when the Loop() returns. func (sync *NodeSync) Loop(logger klog.Logger, done chan struct{}) { ... // 30s触发一次过期 delayTimer := time.NewTimer(timeout) for { select { case op, more := <-sync.opChan: ... // 执行一次同步,更新退避时间(250ms,500ms,...,10s) sync.c.ReportResult(op.run(logger, sync)) // 30s一轮的等待 if !delayTimer.Stop() { <-delayTimer.C } // 30s一轮的等待 case <-delayTimer.C: sync.c.ReportResult((&updateOp{}).run(logger, sync)) } ... // 重置30s过期 delayTimer.Reset(timeout) } } func (op *updateOp) run(logger klog.Logger, sync *NodeSync) error { ... // 获取一下node if op.node == nil { node, err := sync.kubeAPI.Node(ctx, sync.nodeName) ... op.node = node } // 同步node对应cidr aliasRange, err := sync.cloudAlias.Alias(ctx, op.node) ... switch { // 执行分配及更新cloud和node case op.node.Spec.PodCIDR == "" && aliasRange == nil: err = op.allocateRange(ctx, sync, op.node) // 更新node case op.node.Spec.PodCIDR == "" && aliasRange != nil: err = op.updateNodeFromAlias(ctx, sync, op.node, aliasRange) // 更新cloud case op.node.Spec.PodCIDR != "" && aliasRange == nil: err = op.updateAliasFromNode(ctx, sync, op.node) // 校验node cidr和aliasRange合法性 case op.node.Spec.PodCIDR != "" && aliasRange != nil: err = op.validateRange(ctx, sync, op.node, aliasRange) } return err }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
注意
nodesyncer是legacy ipam核心实现,负责基于cluster/cloud模式触发node cidr分配及同步
# 4.rangeAlloc
# 4.1.newlloc
NewCIDRRangeAllocator()会初始化rangeAllocator,以切割cluster cidr为固定大小的多块,向node分配一个或多个cidr。// returns a CIDRAllocator to allocate CIDRs for node. func NewCIDRRangeAllocator(...) (CIDRAllocator, error) { ... // 将cidr划分为cidrset for idx, cidr := range allocatorParams.ClusterCIDRs { cidrSet, err := cidrset.NewCIDRSet(cidr, allocatorParams.NodeCIDRMaskSizes[idx]) if err != nil { return nil, err } cidrSets[idx] = cidrSet } // 实例化rangeAllocator ra := &rangeAllocator{ client: client, clusterCIDRs: allocatorParams.ClusterCIDRs, cidrSets: cidrSets, nodeLister: nodeInformer.Lister(), ... nodeCIDRUpdateChannel: make(chan nodeReservedCIDRs, cidrUpdateQueueSize), ... nodesInProcessing: sets.NewString(), } // cluster cidr和service cidr重复部分标记为占用 if allocatorParams.ServiceCIDR != nil { ra.filterOutServiceRange(logger, allocatorParams.ServiceCIDR) } ... if allocatorParams.SecondaryServiceCIDR != nil { ra.filterOutServiceRange(logger, allocatorParams.SecondaryServiceCIDR) } ... if nodeList != nil { // 先获取node列表 for _, node := range nodeList.Items { if len(node.Spec.PodCIDRs) == 0 { continue } // 标记node.spec.PodCIDR占用 ra.occupyCIDRs(&node) ... } } // node informer回调 nodeInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{ AddFunc: controllerutil.CreateAddNodeHandler(func(node *v1.Node) error { return ra.AllocateOrOccupyCIDR(logger, node) }), UpdateFunc: controllerutil.CreateUpdateNodeHandler(func(_, newNode *v1.Node) error { if len(newNode.Spec.PodCIDRs) == 0 { return ra.AllocateOrOccupyCIDR(logger, newNode) } return nil }), DeleteFunc: controllerutil.CreateDeleteNodeHandler(func(node *v1.Node) error { return ra.ReleaseCIDR(logger, node) }), }) return ra, 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
61
62
63
64
65
66
注意
rangeAlloc基于informer回调分配cidr,分配的node cidr推送到updateChan供worker处理
# 4.2.worker
rangeAllocator分配器内部启动30个worker执行node cidr更新,worker会消费updateChan事件及Patch更新至APIServer。func (r *rangeAllocator) Run(ctx context.Context) { ... // node同步等待 if !cache.WaitForNamedCacheSync("cidrallocator", ctx.Done(), r.nodesSynced) { return } // 激活30个worker for i := 0; i < cidrUpdateWorkers; i++ { go r.worker(ctx) } <-ctx.Done() } func (r *rangeAllocator) worker(ctx context.Context) { for { select { // 监听updateChan case workItem, ok := <-r.nodeCIDRUpdateChannel: ... // 触发更新 if err := r.updateCIDRsAllocation(logger, workItem); err != nil { // 失败重入队 r.nodeCIDRUpdateChannel <- workItem } case <-ctx.Done(): return } } } // updateCIDRsAllocation assigns CIDR to Node and sends an update to the API server. func (r *rangeAllocator) updateCIDRsAllocation(logger klog.Logger, data nodeReservedCIDRs) error { ... // 移除执行进度 defer r.removeNodeFromProcessing(data.nodeName) // cidr列表 cidrsString := ipnetToStringList(data.allocatedCIDRs) // 获取node node, err = r.nodeLister.Get(data.nodeName) ... // 不变性检查 if len(node.Spec.PodCIDRs) == len(data.allocatedCIDRs) { match := true // 对比cidr差异 for idx, cidr := range cidrsString { if node.Spec.PodCIDRs[idx] != cidr { match = false break } } if match { return nil } } // node has cidrs, release the reserved if len(node.Spec.PodCIDRs) != 0 { ... // 释放本次分配的 for idx, cidr := range data.allocatedCIDRs { r.cidrSets[idx].Release(cidr) ... } return nil } // 重试3次设置node cidr for i := 0; i < cidrUpdateRetries; i++ { nodeutil.PatchNodeCIDRs(r.client, types.NodeName(node.Name), cidrsString) ... } ... // 非APIServer连接超时错误 if !apierrors.IsServerTimeout(err) { // 释放本次分配的cidr for idx, cidr := range data.allocatedCIDRs { r.cidrSets[idx].Release(cidr) ... } } return err }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
注意
worker实现相比其它控制器简单一些,基于informer监听回调推送的node cidr更新节点,其它allocator不再赘述,原理类似