fix(master): Leader mismatch during volume cache synchoronize #23117176

Signed-off-by: leonrayang <chl696@sina.com>
This commit is contained in:
leonrayang 2025-03-07 19:52:19 +08:00 committed by zhumingze1108
parent 8d897d1289
commit 27ff3d3280
6 changed files with 26 additions and 28 deletions

View File

@ -16,6 +16,6 @@
"metaNodeReservedMem": "67108864",
"intervalToScanS3Expiration": "180",
"enableLogPanicHook": true,
"enableFollowerCache": false,
"enableFollowerCache": true,
"defaultStorageClass": 1
}

View File

@ -16,6 +16,6 @@
"metaNodeReservedMem": "67108864",
"intervalToScanS3Expiration": "180",
"enableLogPanicHook": true,
"enableFollowerCache": false,
"enableFollowerCache": true,
"defaultStorageClass": 1
}

View File

@ -16,6 +16,6 @@
"metaNodeReservedMem": "67108864",
"intervalToScanS3Expiration": "180",
"enableLogPanicHook": true,
"enableFollowerCache": false,
"enableFollowerCache": true,
"defaultStorageClass": 1
}

View File

@ -235,7 +235,7 @@ func (mgr *followerReadManager) getVolumeDpView() {
}
mgr.rwMutex.Unlock()
if mgr.c.masterClient.Leader() == "" {
if mgr.c.leaderInfo.id == 0 {
log.LogErrorf("followerReadManager.getVolumeDpView but master leader not ready")
return
}
@ -251,6 +251,7 @@ func (mgr *followerReadManager) getVolumeDpView() {
continue
}
mgr.c.masterClient.SetLeader(mgr.c.leaderInfo.addr)
log.LogDebugf("followerReadManager.getVolumeDpView %v leader(%v)", vv.Name, mgr.c.masterClient.Leader())
if view, err = mgr.c.masterClient.ClientAPI().GetDataPartitionsFromLeader(vv.Name); err != nil {
log.LogErrorf("followerReadManager.getVolumeDpView %v GetDataPartitions err %v leader(%v)", vv.Name, err, mgr.c.masterClient.Leader())
@ -270,7 +271,7 @@ func (mgr *followerReadManager) sendFollowerVolumeDpView() {
}
avgSleepTime := time.Second * 5 / time.Duration(len(vols))
for _, vol := range vols {
log.LogDebugf("followerReadManager.getVolumeDpView %v", vol.Name)
log.LogDebugf("followerReadManager.sendFollowerVolumeDpView %v", vol.Name)
if (vol.Status == proto.VolStatusMarkDelete && !vol.Forbidden) || (vol.Status == proto.VolStatusMarkDelete && vol.Forbidden && time.Until(vol.DeleteExecTime) <= 0) {
continue
}
@ -285,6 +286,7 @@ func (mgr *followerReadManager) sendFollowerVolumeDpView() {
if addr == mgr.c.leaderInfo.addr {
continue
}
mgr.c.masterClient.SetLeader(addr)
if err = mgr.c.masterClient.AdminAPI().PutDataPartitions(vol.Name, body); err != nil {
mgr.c.masterClient.SetLeader("")
@ -836,7 +838,6 @@ func (c *Cluster) scheduleToCheckHeartbeat() {
name: "scheduleToCheckHeartbeat_checkDataNodeHeartbeat",
function: func() (fin bool) {
if c.partition != nil && c.partition.IsRaftLeader() {
c.checkLeaderAddr()
c.checkDataNodeHeartbeat()
// update load factor
setOverSoldFactor(c.cfg.ClusterLoadFactor)
@ -881,11 +882,6 @@ func (c *Cluster) scheduleToCheckHeartbeat() {
}()
}
func (c *Cluster) checkLeaderAddr() {
leaderID, _ := c.partition.LeaderTerm()
c.leaderInfo.addr = AddrDatabase[leaderID]
}
func (c *Cluster) checkDataNodeHeartbeat() {
tasks := make([]*proto.AdminTask, 0)
id := uuid.New()

View File

@ -362,26 +362,29 @@ func TestBalanceMetaPartition(t *testing.T) {
}
func TestMasterClientLeaderChange(t *testing.T) {
cluster := &Cluster{
masterClient: masterSDK.NewMasterClient(nil, false),
flashNodeTopo: newFlashNodeTopology(),
}
cluster.t = newTopology()
cluster.BadDataPartitionIds = new(sync.Map)
server := &Server{
cluster: cluster,
leaderInfo: &LeaderInfo{
addr: "",
},
user: &User{},
}
cluster := &Cluster{
masterClient: masterSDK.NewMasterClient(nil, false),
flashNodeTopo: newFlashNodeTopology(),
leaderInfo: server.leaderInfo,
}
server.cluster = cluster
cluster.t = newTopology()
cluster.BadDataPartitionIds = new(sync.Map)
// NOTE: avoid conflict
AddrDatabase[5] = "192.168.0.11:17010"
AddrDatabase[6] = "192.168.0.12:17010"
server.handleLeaderChange(5)
server.handleLeaderChange(6)
masters := cluster.masterClient.GetMasterAddresses()
require.EqualValues(t, 2, len(masters))
require.True(t, cluster.leaderInfo.addr == AddrDatabase[6])
}
func TestCreateVolWithDpCount(t *testing.T) {

View File

@ -27,6 +27,7 @@ import (
// LeaderInfo represents the leader's information
type LeaderInfo struct {
addr string //host:port
id uint64
}
func (m *Server) getCurrAddr() string {
@ -39,30 +40,30 @@ func (m *Server) handleLeaderChange(leader uint64) {
if WarnMetrics != nil {
WarnMetrics.reset()
}
m.leaderInfo.id = 0
m.leaderInfo.addr = ""
return
}
// oldLeaderAddr := m.leaderInfo.addr
m.leaderInfo.addr = AddrDatabase[leader]
log.LogWarnf("action[handleLeaderChange] [%v] ", m.leaderInfo.addr)
m.leaderInfo.id = leader
log.LogWarnf("action[handleLeaderChange] current id [%v] new leader addr [%v] leader id [%v]", m.id, m.leaderInfo.addr, leader)
m.reverseProxy = m.newReverseProxy()
if m.id == leader {
Warn(m.clusterName, fmt.Sprintf("clusterID[%v] leader is changed to %v",
Warn(m.clusterName, fmt.Sprintf("clusterID[%v] current is leader, leader is changed to %v",
m.clusterName, m.leaderInfo.addr))
// if oldLeaderAddr != m.leaderInfo.addr {
m.cluster.checkPersistClusterValue()
m.loadMetadata()
m.cluster.metaReady = true
m.metaReady = true
// }
m.cluster.checkDataNodeHeartbeat()
m.cluster.checkMetaNodeHeartbeat()
m.cluster.checkLcNodeHeartbeat()
m.cluster.lcMgr.startLcScanHandleLeaderChange()
m.cluster.followerReadManager.reSet()
} else {
Warn(m.clusterName, fmt.Sprintf("clusterID[%v] leader is changed to %v",
m.clusterName, m.leaderInfo.addr))
@ -74,8 +75,6 @@ func (m *Server) handleLeaderChange(leader uint64) {
}
m.metaReady = false
m.cluster.metaReady = false
m.cluster.masterClient.AddNode(m.leaderInfo.addr)
m.cluster.masterClient.SetLeader(m.leaderInfo.addr)
if WarnMetrics != nil {
WarnMetrics.reset()
}