Skip to content

Commit b4572f5

Browse files
committed
Read replicas of a master the cluster has declared failed
A replica cut off from its master falls further behind only while somebody writes to the shard. Once the cluster has agreed the master is down, nobody does until a new master is elected, so the replicas' data is the freshest the shard can offer no matter how long ago their links broke. Excluding them after ReplicaLinkDownTolerance turns a stalled failover into a read outage for no gain. The master's state comes from CLUSTER NODES fetched with every reload. Only the agreed flag counts: a node's own suspicion (fail?) does not, and without a readable CLUSTER NODES the tolerance alone decides. A replica that has never synced since it started stays excluded either way.
1 parent 5b6bca9 commit b4572f5

6 files changed

Lines changed: 130 additions & 34 deletions

File tree

‎rediscluster/cluster.go‎

Lines changed: 11 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -117,8 +117,10 @@ type Opts struct {
117117
// WeightProvider - enables to explicitly set weights of replicas (has higher priority than LatencyOrientedRR)
118118
WeightProvider WeightProvider
119119
// ReplicaLinkDownTolerance - a replica whose link with its master is down keeps receiving
120-
// reads while its master_link_down_since_seconds stays below this value. A replica that
121-
// has never synced since it started is never read regardless of it.
120+
// reads while its master_link_down_since_seconds stays below this value. Once the cluster
121+
// declares the master failed, its replicas are read for as long as it stays failed:
122+
// nobody writes to the shard until a new master is elected. A replica that has never
123+
// synced since it started is never read regardless of both.
122124
// default: 60 seconds; negative: a replica with a broken link is never read
123125
ReplicaLinkDownTolerance time.Duration
124126
// Enable connection with TLS
@@ -171,10 +173,11 @@ type clusterConfig struct {
171173
}
172174

173175
type shard struct {
174-
rr uint32
175-
good uint32
176-
addr []string
177-
pingWeights []uint32
176+
rr uint32
177+
good uint32
178+
masterFailed uint32
179+
addr []string
180+
pingWeights []uint32
178181
}
179182
type shardMap map[uint16]*shard
180183
type masterMap map[string]uint16
@@ -389,11 +392,11 @@ func (c *Cluster) control() {
389392
}
390393

391394
func (c *Cluster) reloadMapping() error {
392-
nodes, err := c.slotRangesAndInternalMasterOnly()
395+
nodes, failedMasters, err := c.slotRangesAndInternalMasterOnly()
393396
if err != nil {
394397
return err
395398
}
396-
c.updateMappings(nodes)
399+
c.updateMappings(nodes, failedMasters)
397400
return nil
398401
}
399402

‎rediscluster/cluster_test.go‎

Lines changed: 41 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -518,6 +518,47 @@ func (s *Suite) TestFallbackToSlaveTimeout() {
518518

519519
}
520520

521+
func (s *Suite) TestReadFromReplicaWhileMasterStaysFailed() {
522+
opts := longcheckopts
523+
opts.RoundRobinSeed = alwaysZero{}
524+
opts.ReplicaLinkDownTolerance = -1
525+
cl, err := NewCluster(s.ctx, []string{"127.0.0.1:21100"}, opts)
526+
s.r().Nil(err)
527+
defer cl.Close()
528+
529+
sconn := redis.SyncCtx{cl.WithPolicy(MasterAndSlaves)}
530+
531+
key := slotkey("frozen", s.keys[1], "read")
532+
s.r().Equal("OK", sconn.Do(s.ctx, "SET", key, "1"))
533+
s.waitReplicated(1, 30*time.Second)
534+
535+
s.cl.Node[3].DoSure("CONFIG", "SET", "cluster-replica-no-failover", "yes")
536+
defer s.cl.Node[3].DoSure("CONFIG", "SET", "cluster-replica-no-failover", "no")
537+
s.cl.Node[0].Stop()
538+
defer func() {
539+
s.cl.Node[0].Start()
540+
s.cl.WaitClusterOk()
541+
}()
542+
543+
// The first failed read forces a reload, which marks the replica down.
544+
var res interface{}
545+
for deadline := time.Now().Add(5 * time.Second); time.Now().Before(deadline); time.Sleep(50 * time.Millisecond) {
546+
if res = sconn.Do(s.ctx, "GET", key); redis.AsError(res) != nil {
547+
break
548+
}
549+
}
550+
s.r().True(s.AsError(res).IsOfType(ErrNoAliveConnection), "expected no_alive_connection, got %v", res)
551+
552+
// The testbed runs with cluster-node-timeout of one second, so the cluster agrees the
553+
// master failed within a few seconds and the next reload picks that up.
554+
for deadline := time.Now().Add(15 * time.Second); time.Now().Before(deadline); time.Sleep(200 * time.Millisecond) {
555+
if res = sconn.Do(s.ctx, "GET", key); redis.AsError(res) == nil {
556+
break
557+
}
558+
}
559+
s.Equal([]byte("1"), res)
560+
}
561+
521562
func (s *Suite) TestGetMoved() {
522563
cl, err := NewCluster(s.ctx, []string{"127.0.0.1:21100"}, longcheckopts)
523564
s.r().Nil(err)

‎rediscluster/redisclusterutil/cluster.go‎

Lines changed: 9 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -133,12 +133,14 @@ func parseHostname(metadata []interface{}) string {
133133

134134
// InstanceInfo represents line of CLUSTER NODES result.
135135
type InstanceInfo struct {
136-
Uuid string
137-
Addr string
138-
IP string
139-
Port int
140-
Port2 int
141-
Fail bool
136+
Uuid string
137+
Addr string
138+
IP string
139+
Port int
140+
Port2 int
141+
Fail bool
142+
// PFail: the reporting node suspects the instance (fail?), the cluster has not agreed yet.
143+
PFail bool
142144
MySelf bool
143145
// NoAddr means that node were missed due to misconfiguration.
144146
// More probably, redis instance with other UUID were started on the same port.
@@ -347,6 +349,7 @@ func ParseClusterNodes(res interface{}) (InstanceInfos, error) {
347349
node.Port2, _ = strconv.Atoi(ipp[1])
348350

349351
node.Fail = strings.Contains(parts[2], "fail")
352+
node.PFail = strings.Contains(parts[2], "fail?")
350353
if strings.Contains(parts[2], "slave") {
351354
node.SlaveOf = parts[3]
352355
}

‎rediscluster/redisclusterutil/cluster_test.go‎

Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -269,3 +269,21 @@ func TestParseSlotsInfo_NoHostname(t *testing.T) {
269269

270270
assert.Equal(t, expectedSlots, slots)
271271
}
272+
273+
func TestParseClusterNodes_FailFlags(t *testing.T) {
274+
nodes := "0000000000000000000000000000000000000001 127.0.0.1:21100@31100 master,fail - 0 1789410223036 34 connected 0-5499\n" +
275+
"0000000000000000000000000000000000000002 127.0.0.1:21101@31101 master,fail? - 0 1789410223136 23 connected 5500-10999\n" +
276+
"0000000000000000000000000000000000000003 127.0.0.1:21102@31102 myself,master - 0 0 24 connected 11000-16383\n"
277+
infos, err := ParseClusterNodes([]byte(nodes))
278+
if err != nil {
279+
t.Fatal(err)
280+
}
281+
if len(infos) != 3 {
282+
t.Fatalf("parsed %d nodes, want 3", len(infos))
283+
}
284+
for i, want := range []struct{ fail, pfail bool }{{true, false}, {true, true}, {false, false}} {
285+
if infos[i].Fail != want.fail || infos[i].PFail != want.pfail {
286+
t.Errorf("node %d: Fail=%v PFail=%v, want Fail=%v PFail=%v", i, infos[i].Fail, infos[i].PFail, want.fail, want.pfail)
287+
}
288+
}
289+
}

‎rediscluster/replica_health_internal_test.go‎

Lines changed: 11 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -11,10 +11,11 @@ func TestReplicaHealthy(t *testing.T) {
1111
return []byte("# Replication\r\nrole:slave\r\n" + strings.Join(lines, "\r\n") + "\r\n# Persistence\r\nloading:0\r\n")
1212
}
1313
for _, tc := range []struct {
14-
name string
15-
info []byte
16-
tolerance time.Duration
17-
want bool
14+
name string
15+
info []byte
16+
tolerance time.Duration
17+
masterFailed bool
18+
want bool
1819
}{
1920
{name: "link up", info: info("master_link_status:up"), tolerance: time.Minute, want: true},
2021
{name: "link down for a while", info: info("master_link_status:down", "master_link_down_since_seconds:3"), tolerance: time.Minute, want: true},
@@ -23,10 +24,14 @@ func TestReplicaHealthy(t *testing.T) {
2324
{name: "link down without duration", info: info("master_link_status:down"), tolerance: time.Minute, want: false},
2425
{name: "negative tolerance", info: info("master_link_status:down", "master_link_down_since_seconds:0"), tolerance: -1, want: false},
2526
{name: "loading", info: []byte("# Replication\r\nmaster_link_status:up\r\n# Persistence\r\nloading:1\r\n"), tolerance: time.Minute, want: false},
27+
{name: "master failed, link down beyond tolerance", info: info("master_link_status:down", "master_link_down_since_seconds:3600"), tolerance: time.Minute, masterFailed: true, want: true},
28+
{name: "master failed, negative tolerance", info: info("master_link_status:down", "master_link_down_since_seconds:5"), tolerance: -1, masterFailed: true, want: true},
29+
{name: "master failed, never synced since start", info: info("master_link_status:down", "master_link_down_since_seconds:-1"), tolerance: time.Minute, masterFailed: true, want: false},
30+
{name: "master failed, loading", info: []byte("# Replication\r\nmaster_link_status:down\r\nmaster_link_down_since_seconds:5\r\n# Persistence\r\nloading:1\r\n"), tolerance: time.Minute, masterFailed: true, want: false},
2631
} {
2732
t.Run(tc.name, func(t *testing.T) {
28-
if got := replicaHealthy(tc.info, tc.tolerance); got != tc.want {
29-
t.Errorf("replicaHealthy(%q, %v) = %v, want %v", tc.info, tc.tolerance, got, tc.want)
33+
if got := replicaHealthy(tc.info, tc.tolerance, tc.masterFailed); got != tc.want {
34+
t.Errorf("replicaHealthy(%q, %v, %v) = %v, want %v", tc.info, tc.tolerance, tc.masterFailed, got, tc.want)
3035
}
3136
})
3237
}

‎rediscluster/slotrange.go‎

Lines changed: 40 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -13,17 +13,19 @@ import (
1313

1414
const masterOnlyFlag = 0x4000
1515

16-
func (c *Cluster) slotRangesAndInternalMasterOnly() ([]redisclusterutil.SlotsRange, error) {
16+
func (c *Cluster) slotRangesAndInternalMasterOnly() ([]redisclusterutil.SlotsRange, map[string]struct{}, error) {
1717
nodes := c.getConfig().nodes
1818

1919
var ranges []redisclusterutil.SlotsRange
20+
var failedMasters map[string]struct{}
2021
var err error
2122
Outter:
2223
for _, node := range nodes {
2324
for _, conn := range node.conns {
2425
resp := redis.Sync{conn}.Do("CLUSTER SLOTS")
2526
ranges, err = redisclusterutil.ParseSlotsInfo(resp)
2627
if err == nil {
28+
failedMasters = failedMastersOf(redis.Sync{conn}.Do("CLUSTER NODES"))
2729
break Outter
2830
}
2931
c.report(LogClusterSlotsError{Conn: conn, Error: err})
@@ -32,7 +34,7 @@ Outter:
3234
}
3335
if err != nil {
3436
c.report(LogSlotRangeError{})
35-
return nil, c.err(ErrClusterSlots)
37+
return nil, nil, c.err(ErrClusterSlots)
3638
}
3739

3840
// look for reminder about future migrations
@@ -43,10 +45,28 @@ Outter:
4345
}
4446
c.m.Unlock()
4547

46-
return ranges, nil
48+
return ranges, failedMasters, nil
4749
}
4850

49-
func (c *Cluster) updateMappings(slotRanges []redisclusterutil.SlotsRange) {
51+
// failedMastersOf lists the masters the cluster agrees are down. A node's own
52+
// suspicion (fail?) does not count. Without a readable CLUSTER NODES nothing counts,
53+
// and the link-down tolerance alone decides.
54+
func failedMastersOf(res interface{}) map[string]struct{} {
55+
infos, err := redisclusterutil.ParseClusterNodes(res)
56+
if err != nil {
57+
return nil
58+
}
59+
failed := map[string]struct{}{}
60+
for i := range infos {
61+
ii := &infos[i]
62+
if ii.IsMaster() && ii.Fail && !ii.PFail && ii.HasAddr() {
63+
failed[ii.Addr] = struct{}{}
64+
}
65+
}
66+
return failed
67+
}
68+
69+
func (c *Cluster) updateMappings(slotRanges []redisclusterutil.SlotsRange, failedMasters map[string]struct{}) {
5070
shards := make(map[string][]string)
5171
for _, r := range slotRanges {
5272
shards[r.Addrs[0]] = r.Addrs
@@ -122,19 +142,23 @@ func (c *Cluster) updateMappings(slotRanges []redisclusterutil.SlotsRange) {
122142
return sh
123143
}()
124144

125-
if oldshard != nil {
126-
newConfig.shards[shardno] = oldshard
127-
} else {
128-
shard := &shard{
145+
sh := oldshard
146+
if sh == nil {
147+
sh = &shard{
129148
addr: addrs,
130149
good: (uint32(1) << uint(len(addrs))) - 1,
131150
pingWeights: make([]uint32, len(addrs)),
132151
}
133-
newConfig.shards[shardno] = shard
134-
for i := range shard.pingWeights {
135-
shard.pingWeights[i] = 1
152+
for i := range sh.pingWeights {
153+
sh.pingWeights[i] = 1
136154
}
137155
}
156+
masterFailed := uint32(0)
157+
if _, ok := failedMasters[master]; ok {
158+
masterFailed = 1
159+
}
160+
atomic.StoreUint32(&sh.masterFailed, masterFailed)
161+
newConfig.shards[shardno] = sh
138162
newConfig.masters[addrs[0]] = shardno
139163
random = shardno
140164
}
@@ -230,7 +254,7 @@ func (s *shard) setReplicaInfo(res interface{}, n uint64, tolerance time.Duratio
230254
} else if buf, ok := res.([]byte); !ok {
231255
haserr = true
232256
} else {
233-
haserr = !replicaHealthy(buf, tolerance)
257+
haserr = !replicaHealthy(buf, tolerance, atomic.LoadUint32(&s.masterFailed) != 0)
234258
}
235259
for {
236260
oldstate := atomic.LoadUint32(&s.good)
@@ -252,7 +276,9 @@ func (s *shard) setReplicaInfo(res interface{}, n uint64, tolerance time.Duratio
252276
// replicaHealthy tells whether INFO output describes a replica worth reading from.
253277
// master_link_down_since_seconds is -1 for a replica that has never synced since it
254278
// started, and its dataset is then anything from empty to the RDB it booted from.
255-
func replicaHealthy(info []byte, tolerance time.Duration) bool {
279+
// A replica cut off from a master the cluster has declared failed cannot fall further
280+
// behind: nobody accepts writes for the shard until a new master is elected.
281+
func replicaHealthy(info []byte, tolerance time.Duration, masterFailed bool) bool {
256282
if bytes.Contains(info, []byte("loading:1")) {
257283
return false
258284
}
@@ -263,7 +289,7 @@ func replicaHealthy(info []byte, tolerance time.Duration) bool {
263289
if !ok || since < 0 {
264290
return false
265291
}
266-
return time.Duration(since)*time.Second < tolerance
292+
return masterFailed || time.Duration(since)*time.Second < tolerance
267293
}
268294

269295
func infoInt(info []byte, field string) (int64, bool) {

0 commit comments

Comments
 (0)