Merge pull request #130919 from mengqiy/kubelet_rand_interval

Add random interval to nodeStatusReport interval every time after an actual node status change update or restart
This commit is contained in:
Kubernetes Prow Robot
2025-06-19 15:30:53 -07:00
committed by GitHub
3 changed files with 103 additions and 2 deletions

View File

@@ -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

View File

@@ -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) {

View File

@@ -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)
}
}