From b2db61eb3444826e6c63843ddc2ab8c9fc3b6d09 Mon Sep 17 00:00:00 2001 From: Binh Nguyen Date: Fri, 24 Jul 2026 17:15:40 +0700 Subject: [PATCH] feat(failover): elect the highest-offset replica as master When re-electing a master with no clear existing one (bootstrap or all-down recovery), the operator promoted the oldest pod, which could be a replica that lags behind and discard the data held by more up-to-date replicas. Elect by replication offset instead: query master_repl_offset from each pod and promote the one with the most data, keeping the oldest pod as a deterministic tie-breaker and as the fallback when offsets can't be read. Unreachable pods sort last so they are never elected. Adds redis.Client.GetReplicationOffset (+ mock). Refs upstream spotahome/redis-operator#385. --- metrics/metrics.go | 1 + mocks/service/redis/Client.go | 24 +++++++++++ operator/redisfailover/service/heal.go | 36 +++++++++++++--- operator/redisfailover/service/heal_test.go | 46 +++++++++++++++++++++ service/redis/client.go | 34 +++++++++++++++ 5 files changed, 135 insertions(+), 6 deletions(-) diff --git a/metrics/metrics.go b/metrics/metrics.go index f93ae8888..a22a336e2 100644 --- a/metrics/metrics.go +++ b/metrics/metrics.go @@ -69,6 +69,7 @@ const ( GET_SENTINEL_MONITOR = "SENTINEL_GET_MASTER_INSTANCE" CHECK_SENTINEL_QUORUM = "SENTINEL_CKQUORUM" SLAVE_IS_READY = "CHECK_IF_SLAVE_IS_READY" + GET_REPLICATION_OFFSET = "GET_REPLICATION_OFFSET_OF_INSTANCE" ) var ( // used for grabage collection of metrics diff --git a/mocks/service/redis/Client.go b/mocks/service/redis/Client.go index 6950584ee..272291dab 100644 --- a/mocks/service/redis/Client.go +++ b/mocks/service/redis/Client.go @@ -88,6 +88,30 @@ func (_m *Client) GetSentinelMonitor(ip string) (string, string, error) { return r0, r1, r2 } +// GetReplicationOffset provides a mock function with given fields: ip, port, password +func (_m *Client) GetReplicationOffset(ip string, port string, password string) (int64, error) { + ret := _m.Called(ip, port, password) + + var r0 int64 + var r1 error + if rf, ok := ret.Get(0).(func(string, string, string) (int64, error)); ok { + return rf(ip, port, password) + } + if rf, ok := ret.Get(0).(func(string, string, string) int64); ok { + r0 = rf(ip, port, password) + } else { + r0 = ret.Get(0).(int64) + } + + if rf, ok := ret.Get(1).(func(string, string, string) error); ok { + r1 = rf(ip, port, password) + } else { + r1 = ret.Error(1) + } + + return r0, r1 +} + // GetSlaveOf provides a mock function with given fields: ip, port, password func (_m *Client) GetSlaveOf(ip string, port string, password string) (string, error) { ret := _m.Called(ip, port, password) diff --git a/operator/redisfailover/service/heal.go b/operator/redisfailover/service/heal.go index b47f52543..2d28bedcf 100644 --- a/operator/redisfailover/service/heal.go +++ b/operator/redisfailover/service/heal.go @@ -91,7 +91,21 @@ func (r *RedisFailoverHealer) MakeMaster(ip string, rf *redisfailoverv1.RedisFai return nil } -// SetOldestAsMaster puts all redis to the same master, choosen by order of appearance +// replicationOffsetOrNegative returns the pod's redis replication offset, or -1 +// when it can't be read (e.g. the pod is unreachable) so such a pod sorts last +// and is never preferred as the new master. +func (r *RedisFailoverHealer) replicationOffsetOrNegative(ip, port, password string) int64 { + offset, err := r.redisClient.GetReplicationOffset(ip, port, password) + if err != nil { + return -1 + } + return offset +} + +// SetOldestAsMaster elects a new master and points every other redis at it. The +// candidate with the highest replication offset wins (it has the most data, so +// promoting it loses the least), falling back to the oldest pod as a +// deterministic tie-breaker and when offsets can't be read. func (r *RedisFailoverHealer) SetOldestAsMaster(rf *redisfailoverv1.RedisFailover) error { ssp, err := r.k8sService.GetStatefulSetPods(rf.Namespace, GetRedisName(rf)) if err != nil { @@ -101,17 +115,27 @@ func (r *RedisFailoverHealer) SetOldestAsMaster(rf *redisfailoverv1.RedisFailove return errors.New("number of redis pods are 0") } - // Order the pods so we start by the oldest one - sort.Slice(ssp.Items, func(i, j int) bool { - return ssp.Items[i].CreationTimestamp.Before(&ssp.Items[j].CreationTimestamp) - }) - password, err := k8s.GetRedisPassword(r.k8sService, rf) if err != nil { return err } port := getRedisPort(rf.Spec.Redis.Port) + + // Prefer the most up-to-date replica (highest replication offset). Compute + // each offset once, then order by offset desc, oldest pod as tie-breaker. + offsets := make(map[string]int64, len(ssp.Items)) + for _, pod := range ssp.Items { + offsets[pod.Status.PodIP] = r.replicationOffsetOrNegative(pod.Status.PodIP, port, password) + } + sort.SliceStable(ssp.Items, func(i, j int) bool { + oi, oj := offsets[ssp.Items[i].Status.PodIP], offsets[ssp.Items[j].Status.PodIP] + if oi != oj { + return oi > oj + } + return ssp.Items[i].CreationTimestamp.Before(&ssp.Items[j].CreationTimestamp) + }) + newMasterIP := "" for _, pod := range ssp.Items { if newMasterIP == "" { diff --git a/operator/redisfailover/service/heal_test.go b/operator/redisfailover/service/heal_test.go index 29722dbc8..230c5460f 100644 --- a/operator/redisfailover/service/heal_test.go +++ b/operator/redisfailover/service/heal_test.go @@ -36,6 +36,7 @@ func TestSetOldestAsMasterNewMasterError(t *testing.T) { ms.On("GetStatefulSetPods", namespace, rfservice.GetRedisName(rf)).Once().Return(pods, nil) ms.On("UpdatePodLabels", namespace, mock.AnythingOfType("string"), mock.Anything).Return(nil) mr := &mRedisService.Client{} + mr.On("GetReplicationOffset", mock.Anything, mock.Anything, mock.Anything).Return(int64(0), nil) mr.On("MakeMaster", "0.0.0.0", "0", "").Once().Return(errors.New("")) healer := rfservice.NewRedisFailoverHealer(ms, mr, log.DummyLogger{}) @@ -63,6 +64,7 @@ func TestSetOldestAsMaster(t *testing.T) { ms.On("GetStatefulSetPods", namespace, rfservice.GetRedisName(rf)).Once().Return(pods, nil) ms.On("UpdatePodLabels", namespace, mock.AnythingOfType("string"), mock.Anything).Once().Return(nil) mr := &mRedisService.Client{} + mr.On("GetReplicationOffset", mock.Anything, mock.Anything, mock.Anything).Return(int64(0), nil) mr.On("MakeMaster", "0.0.0.0", "0", "").Once().Return(nil) healer := rfservice.NewRedisFailoverHealer(ms, mr, log.DummyLogger{}) @@ -95,6 +97,7 @@ func TestSetOldestAsMasterMultiplePodsMakeSlaveOfError(t *testing.T) { ms.On("GetStatefulSetPods", namespace, rfservice.GetRedisName(rf)).Once().Return(pods, nil) ms.On("UpdatePodLabels", namespace, mock.AnythingOfType("string"), mock.Anything).Return(nil) mr := &mRedisService.Client{} + mr.On("GetReplicationOffset", mock.Anything, mock.Anything, mock.Anything).Return(int64(0), nil) mr.On("MakeMaster", "0.0.0.0", "0", "").Once().Return(nil) mr.On("MakeSlaveOfWithPort", "1.1.1.1", "0.0.0.0", "0", "").Once().Return(errors.New("")) @@ -128,6 +131,7 @@ func TestSetOldestAsMasterMultiplePods(t *testing.T) { ms.On("GetStatefulSetPods", namespace, rfservice.GetRedisName(rf)).Once().Return(pods, nil) ms.On("UpdatePodLabels", namespace, mock.AnythingOfType("string"), mock.Anything).Return(nil) mr := &mRedisService.Client{} + mr.On("GetReplicationOffset", mock.Anything, mock.Anything, mock.Anything).Return(int64(0), nil) mr.On("MakeMaster", "0.0.0.0", "0", "").Once().Return(nil) mr.On("MakeSlaveOfWithPort", "1.1.1.1", "0.0.0.0", "0", "").Once().Return(nil) @@ -171,6 +175,8 @@ func TestSetOldestAsMasterOrdering(t *testing.T) { ms.On("GetStatefulSetPods", namespace, rfservice.GetRedisName(rf)).Once().Return(pods, nil) ms.On("UpdatePodLabels", namespace, mock.AnythingOfType("string"), mock.Anything).Return(nil) mr := &mRedisService.Client{} + // Equal offsets, so the oldest pod (1.1.1.1) still wins by the tie-breaker. + mr.On("GetReplicationOffset", mock.Anything, mock.Anything, mock.Anything).Return(int64(0), nil) mr.On("MakeMaster", "1.1.1.1", "0", "").Once().Return(nil) mr.On("MakeSlaveOfWithPort", "0.0.0.0", "1.1.1.1", "0", "").Once().Return(nil) @@ -180,6 +186,46 @@ func TestSetOldestAsMasterOrdering(t *testing.T) { assert.NoError(err) } +func TestSetOldestAsMasterPrefersHighestOffset(t *testing.T) { + assert := assert.New(t) + + rf := generateRF() + + // The first/older pod lags behind; the second pod has the higher replication + // offset and must be elected master despite being newer and later in the list. + pods := &corev1.PodList{ + Items: []corev1.Pod{ + { + ObjectMeta: metav1.ObjectMeta{ + CreationTimestamp: metav1.Time{Time: time.Now().Add(-1 * time.Hour)}, + }, + Status: corev1.PodStatus{PodIP: "0.0.0.0"}, + }, + { + ObjectMeta: metav1.ObjectMeta{ + CreationTimestamp: metav1.Time{Time: time.Now()}, + }, + Status: corev1.PodStatus{PodIP: "1.1.1.1"}, + }, + }, + } + + ms := &mK8SService.Services{} + ms.On("GetStatefulSetPods", namespace, rfservice.GetRedisName(rf)).Once().Return(pods, nil) + ms.On("UpdatePodLabels", namespace, mock.AnythingOfType("string"), mock.Anything).Return(nil) + mr := &mRedisService.Client{} + mr.On("GetReplicationOffset", "0.0.0.0", "0", "").Return(int64(100), nil) + mr.On("GetReplicationOffset", "1.1.1.1", "0", "").Return(int64(500), nil) + mr.On("MakeMaster", "1.1.1.1", "0", "").Once().Return(nil) + mr.On("MakeSlaveOfWithPort", "0.0.0.0", "1.1.1.1", "0", "").Once().Return(nil) + + healer := rfservice.NewRedisFailoverHealer(ms, mr, log.DummyLogger{}) + + err := healer.SetOldestAsMaster(rf) + assert.NoError(err) + mr.AssertExpectations(t) +} + func TestSetMasterOnAllMakeMasterError(t *testing.T) { assert := assert.New(t) diff --git a/service/redis/client.go b/service/redis/client.go index e4b09871b..71f5dce09 100644 --- a/service/redis/client.go +++ b/service/redis/client.go @@ -20,6 +20,7 @@ type Client interface { GetNumberSentinelSlavesInMemory(ip string) (int32, error) ResetSentinel(ip string) error GetSlaveOf(ip, port, password string) (string, error) + GetReplicationOffset(ip, port, password string) (int64, error) IsMaster(ip, port, password string) (bool, error) MonitorRedis(ip, monitor, quorum, password string) error MonitorRedisWithPort(ip, monitor, port, quorum, password string) error @@ -49,6 +50,7 @@ const ( slaveNumberREString = "slaves=([0-9]+)" sentinelStatusREString = "status=([a-z]+)" redisMasterHostREString = "master_host:([0-9.]+)" + redisReplOffsetREString = "master_repl_offset:([0-9]+)" redisRoleMaster = "role:master" redisSyncing = "master_sync_in_progress:1" redisMasterSillPending = "master_host:127.0.0.1" @@ -63,6 +65,7 @@ var ( sentinelStatusRE = regexp.MustCompile(sentinelStatusREString) slaveNumberRE = regexp.MustCompile(slaveNumberREString) redisMasterHostRE = regexp.MustCompile(redisMasterHostREString) + redisReplOffsetRE = regexp.MustCompile(redisReplOffsetREString) ) // GetNumberSentinelsInMemory return the number of sentinels that the requested sentinel has @@ -186,6 +189,37 @@ func (c *client) GetSlaveOf(ip, port, password string) (string, error) { return match[1], nil } +// GetReplicationOffset returns the redis master_repl_offset reported by the +// instance's INFO replication. A higher offset means the instance has processed +// more of the replication stream, so it is the most up-to-date candidate to +// promote to master with the least data loss. +func (c *client) GetReplicationOffset(ip, port, password string) (int64, error) { + options := &rediscli.Options{ + Addr: net.JoinHostPort(ip, port), + Password: password, + DB: 0, + } + rClient := rediscli.NewClient(options) + defer func() { _ = rClient.Close() }() + info, err := rClient.Info(context.TODO(), "replication").Result() + if err != nil { + c.metricsRecorder.RecordRedisOperation(metrics.KIND_REDIS, ip, metrics.GET_REPLICATION_OFFSET, metrics.FAIL, getRedisError(err)) + return 0, err + } + match := redisReplOffsetRE.FindStringSubmatch(info) + if len(match) == 0 { + c.metricsRecorder.RecordRedisOperation(metrics.KIND_REDIS, ip, metrics.GET_REPLICATION_OFFSET, metrics.SUCCESS, metrics.NOT_APPLICABLE) + return 0, nil + } + offset, err := strconv.ParseInt(match[1], 10, 64) + if err != nil { + c.metricsRecorder.RecordRedisOperation(metrics.KIND_REDIS, ip, metrics.GET_REPLICATION_OFFSET, metrics.FAIL, getRedisError(err)) + return 0, err + } + c.metricsRecorder.RecordRedisOperation(metrics.KIND_REDIS, ip, metrics.GET_REPLICATION_OFFSET, metrics.SUCCESS, metrics.NOT_APPLICABLE) + return offset, nil +} + func (c *client) IsMaster(ip, port, password string) (bool, error) { options := &rediscli.Options{ Addr: net.JoinHostPort(ip, port),