fix(master): fix deadlock for progress of WarnMissingDp

Signed-off-by: chihe <chihe@oppo.com>
This commit is contained in:
chihe 2024-06-14 15:59:25 +08:00 committed by longerfly
parent e9804d0043
commit 3325185e75

View File

@ -253,6 +253,8 @@ func (m *warningMetrics) reset() {
// The caller is responsible for lock
func (m *warningMetrics) deleteMissingDp(missingDpAddrSet addrSet, clusterName, dpId, addr string) {
m.dpMissingReplicaMutex.Lock()
defer m.dpMissingReplicaMutex.Unlock()
if len(missingDpAddrSet.addrs) == 0 {
return
}
@ -260,8 +262,6 @@ func (m *warningMetrics) deleteMissingDp(missingDpAddrSet addrSet, clusterName,
if _, ok := missingDpAddrSet.addrs[addr]; !ok {
return
}
m.dpMissingReplicaMutex.Lock()
defer m.dpMissingReplicaMutex.Unlock()
replicaAlive := m.dpMissingReplicaInfo[dpId].replicaAlive
replicaNum := m.dpMissingReplicaInfo[dpId].replicaNum
@ -278,13 +278,13 @@ func (m *warningMetrics) deleteMissingDp(missingDpAddrSet addrSet, clusterName,
// leader only
func (m *warningMetrics) WarnMissingDp(clusterName, addr string, partitionID uint64, report bool) {
m.dpMissingReplicaMutex.Lock()
defer m.dpMissingReplicaMutex.Unlock()
if clusterName != m.cluster.Name {
return
}
m.dpMissingReplicaMutex.Lock()
id := strconv.FormatUint(partitionID, 10)
if !report {
m.dpMissingReplicaMutex.Unlock()
m.deleteMissingDp(m.dpMissingReplicaInfo[id], clusterName, id, addr)
return
}
@ -295,6 +295,7 @@ func (m *warningMetrics) WarnMissingDp(clusterName, addr string, partitionID uin
// m.dpMissingReplicaInfo[id].addrs = make(addrSet)
}
m.dpMissingReplicaInfo[id].addrs[addr] = voidVal
m.dpMissingReplicaMutex.Unlock()
}
// leader only