fix(datanode): reduce the locking time when building heart beat

Signed-off-by: chihe <chihe@oppo.com>
This commit is contained in:
chihe 2024-07-18 10:31:25 +08:00 committed by AmazingChi
parent a2da2149c3
commit 8d68de0acc
6 changed files with 87 additions and 56 deletions

View File

@ -390,7 +390,6 @@ func newDataPartition(dpCfg *dataPartitionCfg, disk *Disk, isCreate bool) (dp *D
dp = partition
go partition.statusUpdateScheduler()
go partition.startEvict()
go partition.validatePeers()
if isCreate {
if err = dp.getVerListFromMaster(); err != nil {
log.LogErrorf("action[newDataPartition] vol %v dp %v loadFromMaster verList failed err %v", dp.volumeID, dp.partitionID, err)
@ -775,6 +774,7 @@ func (dp *DataPartition) PersistMetadata() (err error) {
func (dp *DataPartition) statusUpdateScheduler() {
ticker := time.NewTicker(time.Minute)
snapshotTicker := time.NewTicker(time.Minute * 5)
peersTicker := time.NewTicker(10 * time.Second)
var index int
for {
select {
@ -798,6 +798,8 @@ func (dp *DataPartition) statusUpdateScheduler() {
}
case <-snapshotTicker.C:
dp.ReloadSnapshot()
case <-peersTicker.C:
dp.validatePeers()
case <-dp.stopC:
ticker.Stop()
snapshotTicker.Stop()
@ -1529,42 +1531,37 @@ func (dp *DataPartition) info() string {
}
func (dp *DataPartition) validatePeers() {
ticker := time.NewTicker(10 * time.Second)
for {
select {
case <-ticker.C:
dataNodes := dp.dataNode.space.getDataNodeIDs()
for _, peer := range dp.config.Peers {
for _, dn := range dataNodes {
if dn.Addr == peer.Addr && dn.ID != peer.ID {
log.LogWarnf("dp %v find expired peer %v(expected %v_%v)", dp.info(), peer, dn.ID, dn.Addr)
newReq := &proto.RemoveDataPartitionRaftMemberRequest{
PartitionId: dp.partitionID,
Force: true,
RemovePeer: peer,
}
reqData, err := json.Marshal(newReq)
if err != nil {
log.LogWarnf("dp %v marshal newReq %v failed %v", dp.info(), newReq, err)
continue
}
cc := &raftProto.ConfChange{
Type: raftProto.ConfRemoveNode,
Peer: raftProto.Peer{
ID: peer.ID,
},
Context: reqData,
}
dp.dataNode.space.raftStore.RaftServer().RemoveRaftForce(dp.partitionID, cc)
dp.ApplyMemberChange(cc, 0)
dp.PersistMetadata()
log.LogWarnf("dp %v remove expired peer %v", dp.info(), peer)
}
dataNodes := dp.dataNode.space.getDataNodeIDs()
for _, peer := range dp.config.Peers {
for _, dn := range dataNodes {
if dn.Addr == peer.Addr && dn.ID != peer.ID {
log.LogWarnf("dp %v find expired peer %v(expected %v_%v)", dp.info(), peer, dn.ID, dn.Addr)
newReq := &proto.RemoveDataPartitionRaftMemberRequest{
PartitionId: dp.partitionID,
Force: true,
RemovePeer: peer,
}
reqData, err := json.Marshal(newReq)
if err != nil {
log.LogWarnf("dp %v marshal newReq %v failed %v", dp.info(), newReq, err)
continue
}
cc := &raftProto.ConfChange{
Type: raftProto.ConfRemoveNode,
Peer: raftProto.Peer{
ID: peer.ID,
},
Context: reqData,
}
dp.dataNode.space.raftStore.RaftServer().RemoveRaftForce(dp.partitionID, cc)
dp.ApplyMemberChange(cc, 0)
dp.PersistMetadata()
log.LogWarnf("dp %v remove expired peer %v", dp.info(), peer)
}
case <-dp.stopC:
ticker.Stop()
return
}
}
}
func (dp *DataPartition) GetExtentCountWithoutLock() int {
return dp.extentStore.GetExtentCountWithoutLock()
}

View File

@ -108,7 +108,7 @@ func (s *DataNode) getRaftStatus(w http.ResponseWriter, r *http.Request) {
func (s *DataNode) getPartitionsAPI(w http.ResponseWriter, r *http.Request) {
partitions := make([]interface{}, 0)
s.space.RangePartitions(func(dp *DataPartition) bool {
s.space.RangePartitions(func(dp *DataPartition, testID string) bool {
partition := &struct {
ID uint64 `json:"id"`
Size int `json:"size"`
@ -408,7 +408,7 @@ func (s *DataNode) getMetricsDegrade(w http.ResponseWriter, r *http.Request) {
func (s *DataNode) genClusterVersionFile(w http.ResponseWriter, r *http.Request) {
paths := make([]string, 0)
s.space.RangePartitions(func(partition *DataPartition) bool {
s.space.RangePartitions(func(partition *DataPartition, testID string) bool {
paths = append(paths, partition.disk.Path)
return true
}, "")

View File

@ -33,6 +33,7 @@ import (
"github.com/cubefs/cubefs/util/loadutil"
"github.com/cubefs/cubefs/util/log"
"github.com/cubefs/cubefs/util/strutil"
"github.com/google/uuid"
"github.com/shirou/gopsutil/disk"
)
@ -245,7 +246,7 @@ func (manager *SpaceManager) GetRaftStore() (raftStore raftstore.RaftStore) {
return manager.raftStore
}
func (manager *SpaceManager) RangePartitions(f func(partition *DataPartition) bool, reqID string) {
func (manager *SpaceManager) RangePartitions(f func(partition *DataPartition, testID string) bool, reqID string) {
if f == nil {
return
}
@ -256,7 +257,23 @@ func (manager *SpaceManager) RangePartitions(f func(partition *DataPartition) bo
partitions = append(partitions, dp)
}
manager.partitionMutex.RUnlock()
log.LogDebugf("RangePartitions req(%v) get lock cost %v", reqID, time.Now().Sub(begin))
testID := uuid.New().String()
log.LogDebugf("RangePartitions req(%v) get lock cost %v testID %v", reqID, time.Now().Sub(begin), testID)
//for _, partition := range partitions {
// begin2 := time.Now()
// if !f(partition, testID) {
// break
// }
// interval := time.Now().Sub(begin2)
// interval2 := time.Now().Sub(begin)
// if interval > time.Millisecond {
// log.LogDebugf("RangePartitions req(%v) execute fun for dp %v cost %v testID %v too long goroutine %v cost from begin %v",
// reqID, partition.partitionID, interval, testID, runtime.NumGoroutine(), interval2)
// }
// log.LogDebugf("RangePartitions req(%v) execute fun cost %v testID %v", reqID, interval, testID)
//
//}
var wg sync.WaitGroup
partitionsCh := make(chan *DataPartition)
@ -266,7 +283,7 @@ func (manager *SpaceManager) RangePartitions(f func(partition *DataPartition) bo
go func() {
defer wg.Done()
for partition := range partitionsCh {
if !f(partition) {
if !f(partition, testID) {
break
}
}
@ -277,7 +294,7 @@ func (manager *SpaceManager) RangePartitions(f func(partition *DataPartition) bo
}
close(partitionsCh)
wg.Wait()
log.LogDebugf("RangePartitions req(%v) traverse dps cost %v", reqID, time.Now().Sub(begin))
log.LogDebugf("RangePartitions req(%v) traverse dps %v cost %v testID %v", reqID, len(partitions), time.Now().Sub(begin), testID)
}
func (manager *SpaceManager) GetDisks() (disks []*Disk) {
@ -542,8 +559,6 @@ func (manager *SpaceManager) DetachDataPartition(partitionID uint64) {
}
func (manager *SpaceManager) CreatePartition(request *proto.CreateDataPartitionRequest) (dp *DataPartition, err error) {
manager.partitionMutex.Lock()
defer manager.partitionMutex.Unlock()
dpCfg := &dataPartitionCfg{
PartitionID: request.PartitionId,
VolName: request.VolumeId,
@ -561,7 +576,7 @@ func (manager *SpaceManager) CreatePartition(request *proto.CreateDataPartitionR
}
log.LogInfof("action[CreatePartition] dp %v dpCfg.Peers %v request.Members %v",
dpCfg.PartitionID, dpCfg.Peers, request.Members)
dp = manager.partitions[dpCfg.PartitionID]
dp = manager.Partition(dpCfg.PartitionID)
if dp != nil {
if err = dp.IsEqualCreateDataPartitionRequest(request); err != nil {
return nil, err
@ -576,7 +591,9 @@ func (manager *SpaceManager) CreatePartition(request *proto.CreateDataPartitionR
if dp, err = CreateDataPartition(dpCfg, disk, request); err != nil {
return
}
manager.partitionMutex.Lock()
manager.partitions[dp.partitionID] = dp
manager.partitionMutex.Unlock()
return
}
@ -623,8 +640,10 @@ func (s *DataNode) buildHeartBeatResponse(response *proto.DataNodeHeartbeatRespo
response.PartitionReports = make([]*proto.DataPartitionReport, 0)
space := s.space
begin := time.Now()
space.RangePartitions(func(partition *DataPartition) bool {
var respLock sync.Mutex
space.RangePartitions(func(partition *DataPartition, testID string) bool {
leaderAddr, isLeader := partition.IsRaftLeader()
begin2 := time.Now()
vr := &proto.DataPartitionReport{
VolName: partition.volumeID,
PartitionID: uint64(partition.partitionID),
@ -633,16 +652,20 @@ func (s *DataNode) buildHeartBeatResponse(response *proto.DataNodeHeartbeatRespo
Used: uint64(partition.Used()),
DiskPath: partition.Disk().Path,
IsLeader: isLeader,
ExtentCount: partition.GetExtentCount(),
ExtentCount: partition.GetExtentCountWithoutLock(),
NeedCompare: true,
DecommissionRepairProgress: partition.decommissionRepairProgress,
LocalPeers: partition.config.Peers,
TriggerDiskError: atomic.LoadUint64(&partition.diskErrCnt) > 0,
}
log.LogDebugf("action[Heartbeats] dpid(%v), status(%v) total(%v) used(%v) leader(%v) isLeader(%v) TriggerDiskError(%v).",
vr.PartitionID, vr.PartitionStatus, vr.Total, vr.Used, leaderAddr, vr.IsLeader, vr.TriggerDiskError)
log.LogDebugf("action[Heartbeats] dpid(%v), status(%v) total(%v) used(%v) leader(%v) isLeader(%v) "+
"TriggerDiskError(%v) reqId(%v) testID(%v)cost(%v).",
vr.PartitionID, vr.PartitionStatus, vr.Total, vr.Used, leaderAddr, vr.IsLeader, vr.TriggerDiskError,
reqID, testID, time.Now().Sub(begin2))
respLock.Lock()
response.PartitionReports = append(response.PartitionReports, vr)
respLock.Unlock()
begin2 = time.Now()
if len(volNames) != 0 {
if _, ok := volNames[partition.volumeID]; ok {
partition.SetForbidden(true)
@ -657,8 +680,8 @@ func (s *DataNode) buildHeartBeatResponse(response *proto.DataNodeHeartbeatRespo
size = proto.DefaultDpRepairBlockSize
}
}
log.LogDebugf("action[Heartbeats] volume(%v) dp(%v) repair block size(%v) current size(%v)",
partition.volumeID, partition.partitionID, size, partition.GetRepairBlockSize())
log.LogDebugf("action[Heartbeats] volume(%v) dp(%v) repair block size(%v) current size(%v) reqId(%v) testID(testID) cost(%v)",
partition.volumeID, partition.partitionID, size, partition.GetRepairBlockSize(), reqID, testID, time.Now().Sub(begin2))
if partition.GetRepairBlockSize() != size {
partition.SetRepairBlockSize(size)
}

View File

@ -321,9 +321,13 @@ func (s *DataNode) commitCreateVersion(req *proto.MultiVersionOpRequest) (err er
}
s.space.partitionMutex.RLock()
defer s.space.partitionMutex.RUnlock()
resultCh := make(chan error, len(s.space.partitions))
for _, partition := range s.space.partitions {
partitions := make([]*DataPartition, 0)
for _, dp := range s.space.partitions {
partitions = append(partitions, dp)
}
s.space.partitionMutex.RUnlock()
resultCh := make(chan error, len(partitions))
for _, partition := range partitions {
if partition.config.VolName != req.VolumeID {
continue
}
@ -1972,7 +1976,7 @@ func (s *DataNode) handlePacketToQueryBadDiskRecoverProgress(p *repl.Packet) {
log.LogWarnf("action[handlePacketToRecoverBadDisk] disk(%v) is not found err(%v).", request.DiskPath, err)
return
}
total := disk.space.getPartitionIds()
total := disk.DataPartitionList()
badDpList := disk.GetDiskErrPartitionList()
resp := &proto.BadDiskRecoverProgress{
TotalPartitionsNum: len(total),

View File

@ -1153,7 +1153,9 @@ directly:
// return errors.NewErrorf("set RestoreReplicaMetaForbidden failed")
// }
// wait for checkReplicaMeta ended
time.Sleep(3 * time.Second)
log.LogWarnf("action[MarkDecommissionStatus] dp [%d]wait for setting restore replica forbidden",
partition.PartitionID)
time.Sleep(1 * time.Second)
continue
}
break
@ -1172,6 +1174,7 @@ directly:
log.LogWarnf("action[MarkDecommissionStatus] dp [%d]wait for setting restore replica forbidden",
partition.PartitionID)
time.Sleep(1 * time.Second)
continue
}
break
}

View File

@ -1770,3 +1770,7 @@ func (s *ExtentStore) ExtentBatchUnlockNormalExtent(ext []*proto.ExtentKey) {
s.extentLockMap = make(map[uint64]proto.GcFlag)
s.extentLock = false
}
func (s *ExtentStore) GetExtentCountWithoutLock() (count int) {
return len(s.extentInfoMap)
}