diff --git a/datanode/partition.go b/datanode/partition.go index cee665b2e..6a3040949 100644 --- a/datanode/partition.go +++ b/datanode/partition.go @@ -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() +} diff --git a/datanode/server_handler.go b/datanode/server_handler.go index a3c197d46..ee34521bb 100644 --- a/datanode/server_handler.go +++ b/datanode/server_handler.go @@ -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 }, "") diff --git a/datanode/space_manager.go b/datanode/space_manager.go index efc7fb131..5eed6a4b1 100644 --- a/datanode/space_manager.go +++ b/datanode/space_manager.go @@ -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) } diff --git a/datanode/wrap_operator.go b/datanode/wrap_operator.go index e9362989f..ce523cbef 100644 --- a/datanode/wrap_operator.go +++ b/datanode/wrap_operator.go @@ -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), diff --git a/master/data_partition.go b/master/data_partition.go index 9e0fad3df..425628c4e 100644 --- a/master/data_partition.go +++ b/master/data_partition.go @@ -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 } diff --git a/storage/extent_store.go b/storage/extent_store.go index 6bc7b664d..475ca6719 100644 --- a/storage/extent_store.go +++ b/storage/extent_store.go @@ -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) +}