diff --git a/pkg/kubelet/kubelet.go b/pkg/kubelet/kubelet.go index 61f27a2aaf1..7bd32f1b3f4 100644 --- a/pkg/kubelet/kubelet.go +++ b/pkg/kubelet/kubelet.go @@ -1277,6 +1277,12 @@ type Kubelet struct { // status to master. It is only used when node lease feature is enabled. nodeStatusReportFrequency time.Duration + // delayAfterNodeStatusChange is the one-time random duration that we add to the next node status report interval + // every time when there's an actual node status change or kubelet restart. But all future node status update that + // is not caused by real status change will stick with nodeStatusReportFrequency. The random duration is a uniform + // distribution over [-0.5*nodeStatusReportFrequency, 0.5*nodeStatusReportFrequency] + delayAfterNodeStatusChange time.Duration + // lastStatusReportTime is the time when node status was last reported. lastStatusReportTime time.Time diff --git a/pkg/kubelet/kubelet_node_status.go b/pkg/kubelet/kubelet_node_status.go index c562f3501be..9abbc1525e6 100644 --- a/pkg/kubelet/kubelet_node_status.go +++ b/pkg/kubelet/kubelet_node_status.go @@ -19,6 +19,7 @@ package kubelet import ( "context" "fmt" + "math/rand" "net" goruntime "runtime" "sort" @@ -516,13 +517,26 @@ func (kl *Kubelet) tryUpdateNodeStatus(ctx context.Context, tryNumber int) error } node, changed := kl.updateNode(ctx, originalNode) - shouldPatchNodeStatus := changed || kl.clock.Since(kl.lastStatusReportTime) >= kl.nodeStatusReportFrequency + shouldPatchNodeStatus := changed || kl.isUpdateStatusPeriodExpired() if !shouldPatchNodeStatus { kl.markVolumesFromNode(node) return nil } + // There are 3 possible conditions that make shouldPatchNodeStatus to be true: + // 1. node is changed + // 2. isUpdateStatusPeriodExpired returns true due to lastStatusReportTime has Zero value. This will happen when kubelet restarts. + // 3. isUpdateStatusPeriodExpired returns true due to lastStatusReportTime expires with non-zero value. + // We want to calculate a new random delay for condition 1 and 2, so that we can avoid all the periodic node status + // updates to reach the apiserver at the same time. + // When condition 3 happens, random interval has already been used, and we want to reset the random delay, so that + // the node updates its status with fixed interval going forward. + if changed || kl.lastStatusReportTime.IsZero() { + kl.delayAfterNodeStatusChange = kl.calculateDelay() + } else { + kl.delayAfterNodeStatusChange = 0 + } updatedNode, err := kl.patchNodeStatus(originalNode, node) if err == nil { kl.markVolumesFromNode(updatedNode) @@ -530,6 +544,14 @@ func (kl *Kubelet) tryUpdateNodeStatus(ctx context.Context, tryNumber int) error return err } +func (kl *Kubelet) isUpdateStatusPeriodExpired() bool { + return kl.clock.Since(kl.lastStatusReportTime) >= kl.nodeStatusReportFrequency+kl.delayAfterNodeStatusChange +} + +func (kl *Kubelet) calculateDelay() time.Duration { + return time.Duration(float64(kl.nodeStatusReportFrequency) * (-0.5 + rand.Float64())) +} + // updateNode creates a copy of originalNode and runs update logic on it. // It returns the updated node object and a bool indicating if anything has been changed. func (kl *Kubelet) updateNode(ctx context.Context, originalNode *v1.Node) (*v1.Node, bool) { diff --git a/pkg/kubelet/kubelet_node_status_test.go b/pkg/kubelet/kubelet_node_status_test.go index 2a98e74a9b3..a6f587d55ea 100644 --- a/pkg/kubelet/kubelet_node_status_test.go +++ b/pkg/kubelet/kubelet_node_status_test.go @@ -849,7 +849,10 @@ func TestUpdateNodeStatusWithLease(t *testing.T) { // Since this test retroactively overrides the stub container manager, // we have to regenerate default status setters. kubelet.setNodeStatusFuncs = kubelet.defaultNodeStatusFuncs() - kubelet.nodeStatusReportFrequency = time.Minute + // You will add up to 50% of nodeStatusReportFrequency of additional random latency for + // kubelet to determine if update node status is needed due to time passage. We need to + // take that into consideration to ensure this test pass all time. + kubelet.nodeStatusReportFrequency = 30 * time.Second kubeClient := testKubelet.fakeKubeClient existingNode := &v1.Node{ObjectMeta: metav1.ObjectMeta{Name: testKubeletHostname}} @@ -3088,3 +3091,73 @@ func TestUpdateNodeAddresses(t *testing.T) { }) } } + +func TestIsUpdateStatusPeriodExpired(t *testing.T) { + testcases := []struct { + name string + lastStatusReportTime time.Time + delayAfterNodeStatusChange time.Duration + expectExpired bool + }{ + { + name: "no status update before and no delay", + lastStatusReportTime: time.Time{}, + delayAfterNodeStatusChange: 0, + expectExpired: true, + }, + { + name: "no status update before and existing delay", + lastStatusReportTime: time.Time{}, + delayAfterNodeStatusChange: 30 * time.Second, + expectExpired: true, + }, + { + name: "not expired and no delay", + lastStatusReportTime: time.Now().Add(-4 * time.Minute), + delayAfterNodeStatusChange: 0, + expectExpired: false, + }, + { + name: "not expired", + lastStatusReportTime: time.Now().Add(-5 * time.Minute), + delayAfterNodeStatusChange: time.Minute, + expectExpired: false, + }, + { + name: "expired", + lastStatusReportTime: time.Now().Add(-4 * time.Minute), + delayAfterNodeStatusChange: -2 * time.Minute, + expectExpired: true, + }, + { + name: "Delay exactly at threshold", + lastStatusReportTime: time.Now().Add(-5 * time.Minute), + delayAfterNodeStatusChange: 0, + expectExpired: true, + }, + } + + testKubelet := newTestKubelet(t, false /* controllerAttachDetachEnabled */) + defer testKubelet.Cleanup() + kubelet := testKubelet.kubelet + kubelet.nodeStatusReportFrequency = 5 * time.Minute + + for _, tc := range testcases { + kubelet.lastStatusReportTime = tc.lastStatusReportTime + kubelet.delayAfterNodeStatusChange = tc.delayAfterNodeStatusChange + expired := kubelet.isUpdateStatusPeriodExpired() + assert.Equal(t, tc.expectExpired, expired, tc.name) + } +} + +func TestCalculateDelay(t *testing.T) { + testKubelet := newTestKubelet(t, false /* controllerAttachDetachEnabled */) + defer testKubelet.Cleanup() + kubelet := testKubelet.kubelet + kubelet.nodeStatusReportFrequency = 5 * time.Minute + + for i := 0; i < 100; i++ { + randomDelay := kubelet.calculateDelay() + assert.LessOrEqual(t, randomDelay.Abs(), kubelet.nodeStatusReportFrequency/2) + } +}