fix(datanode): optimize the process of constructing heartbeat packets.

Signed-off-by: chihe <chihe@oppo.com>
This commit is contained in:
chihe 2024-07-05 18:54:00 +08:00 committed by AmazingChi
parent 8ad898a1e4
commit 7d05d431e6
4 changed files with 35 additions and 13 deletions

View File

@ -78,7 +78,10 @@ func (s *DataNode) getDiskAPI(w http.ResponseWriter, r *http.Request) {
func (s *DataNode) getStatAPI(w http.ResponseWriter, r *http.Request) {
response := &proto.DataNodeHeartbeatResponse{}
s.buildHeartBeatResponse(response)
forbiddenVols := make(map[string]struct{})
volDpRepairBlockSizes := make(map[string]uint64)
s.buildHeartBeatResponse(response, forbiddenVols, volDpRepairBlockSizes)
s.buildSuccessResp(w, response)
}

View File

@ -585,7 +585,8 @@ func (manager *SpaceManager) DeletePartition(dpID uint64, force bool) (err error
return nil
}
func (s *DataNode) buildHeartBeatResponse(response *proto.DataNodeHeartbeatResponse) {
func (s *DataNode) buildHeartBeatResponse(response *proto.DataNodeHeartbeatResponse,
volNames map[string]struct{}, dpRepairBlockSize map[string]uint64) {
response.Status = proto.TaskSucceeds
stat := s.space.Stats()
stat.Lock()
@ -623,6 +624,26 @@ func (s *DataNode) buildHeartBeatResponse(response *proto.DataNodeHeartbeatRespo
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)
response.PartitionReports = append(response.PartitionReports, vr)
if len(volNames) != 0 {
if _, ok := volNames[partition.volumeID]; ok {
partition.SetForbidden(true)
} else {
partition.SetForbidden(false)
}
}
size := uint64(proto.DefaultDpRepairBlockSize)
if len(dpRepairBlockSize) != 0 {
var ok bool
if size, ok = dpRepairBlockSize[partition.volumeID]; !ok {
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())
if partition.GetRepairBlockSize() != size {
partition.SetRepairBlockSize(size)
}
return true
})

View File

@ -559,18 +559,16 @@ func (s *DataNode) handleHeartbeatPacket(p *repl.Packet) {
log.LogDebugf("handleHeartbeatPacket checkDecommissionDisks req(%v) cost %v",
task.RequestID, time.Now().Sub(begin))
s.buildHeartBeatResponse(response)
forbiddenVols := make(map[string]struct{})
for _, vol := range request.ForbiddenVols {
if _, ok := forbiddenVols[vol]; !ok {
forbiddenVols[vol] = struct{}{}
}
}
s.buildHeartBeatResponse(response, forbiddenVols, request.VolDpRepairBlockSize)
log.LogDebugf("handleHeartbeatPacket buildHeartBeatResponse req(%v) cost %v",
task.RequestID, time.Now().Sub(begin))
// set volume forbidden
s.checkVolumeForbidden(request.ForbiddenVols)
log.LogDebugf("handleHeartbeatPacket checkVolumeForbidden req(%v) cost %v",
task.RequestID, time.Now().Sub(begin))
s.diskQosEnableFromMaster = request.EnableDiskQos
s.checkVolumeDpRepairBlockSize(request.VolDpRepairBlockSize)
log.LogDebugf("handleHeartbeatPacket checkVolumeDpRepairBlockSize req(%v) cost %v",
task.RequestID, time.Now().Sub(begin))
var needUpdate bool
for _, pair := range []struct {
replace uint64

View File

@ -7315,8 +7315,8 @@ func (m *Server) recoverBadDisk(w http.ResponseWriter, r *http.Request) {
}
}
if !found {
sendErrReply(w, r, &proto.HTTPReply{Code: proto.ErrCodeParamError, Msg: fmt.Sprintf("disk %v is not reported"+
"from dataNode %v", diskPath, offLineAddr)})
sendErrReply(w, r, &proto.HTTPReply{Code: proto.ErrCodeParamError, Msg: fmt.Sprintf("disk %v is not bad "+
"disk on dataNode %v, do not support recover", diskPath, offLineAddr)})
return
}
err = dataNode.createTaskToRecoverBadDisk(diskPath)