diff --git a/datanode/disk.go b/datanode/disk.go index 780f013d6..d5961aee1 100644 --- a/datanode/disk.go +++ b/datanode/disk.go @@ -155,6 +155,22 @@ func NewDisk(path string, reservedSpace, diskRdonlySpace uint64, maxErrCnt int, return } +func NewBrokenDisk(path string, reservedSpace, diskRdonlySpace uint64, maxErrCnt int, space *SpaceManager, diskEnableReadRepairExtentLimit bool) (d *Disk) { + d = &Disk{ + Path: path, + ReservedSpace: reservedSpace, + MaxErrCnt: maxErrCnt, + Status: proto.Unavailable, + RejectWrite: true, + space: space, + dataNode: space.dataNode, + partitionMap: make(map[uint64]*DataPartition), + DiskErrPartitionSet: sync.Map{}, + enableExtentRepairReadLimit: diskEnableReadRepairExtentLimit, + } + return +} + func (d *Disk) MarkDecommissionStatus(decommission bool) { probePath := path.Join(d.Path, DecommissionDiskMark) var err error @@ -199,7 +215,16 @@ func (d *Disk) GetDiskPartition() *disk.PartitionStat { return d.diskPartition } +func (d *Disk) isBrokenDisk() (ok bool) { + ok = d.Status == proto.Unavailable && d.limitRead == nil && d.limitWrite == nil + return +} + func (d *Disk) updateQosLimiter() { + if d.isBrokenDisk() { + log.LogInfof("[updateQosLimiter] disk(%v) is broken", d.Path) + return + } if d.dataNode.diskReadFlow > 0 { d.limitFactor[proto.FlowReadType].SetLimit(rate.Limit(d.dataNode.diskReadFlow)) } diff --git a/datanode/server.go b/datanode/server.go index e9119221a..dd7ead065 100644 --- a/datanode/server.go +++ b/datanode/server.go @@ -431,6 +431,30 @@ func (s *DataNode) newSpaceManager(cfg *config.Config) (err error) { return } +func (s *DataNode) getBrokenDisks() (disks map[string]interface{}, err error) { + var dataNode *proto.DataNodeInfo + for i := 0; i < 3; i++ { + dataNode, err = MasterClient.NodeAPI().GetDataNode(s.localServerAddr) + if err != nil { + log.LogErrorf("action[getBrokenDisks]: getDataNode error %v", err) + continue + } + disks = make(map[string]interface{}) + break + } + + if disks == nil { + log.LogErrorf("action[getBrokenDisks]: failed to get datanode(%v), err(%v)", s.localServerAddr, err) + err = fmt.Errorf("failed to get datanode %v", s.localServerAddr) + return + } + log.LogInfof("[getBrokenDisks] data node(%v) broken disks(%v)", dataNode.Addr, dataNode.BadDisks) + for _, disk := range dataNode.BadDisks { + disks[disk] = 1 + } + return +} + func (s *DataNode) startSpaceManager(cfg *config.Config) (err error) { diskRdonlySpace := uint64(cfg.GetInt64(CfgDiskRdonlySpace)) if diskRdonlySpace < DefaultDiskRetainMin { @@ -453,6 +477,13 @@ func (s *DataNode) startSpaceManager(cfg *config.Config) (err error) { } } + brokenDisks, err := s.getBrokenDisks() + if err != nil { + log.LogErrorf("[startSpaceManager] failed to get broken disks, err(%v)", err) + return + } + log.LogInfof("[startSpaceManager] broken disks(%v)", brokenDisks) + var wg sync.WaitGroup for _, d := range paths { log.LogDebugf("action[startSpaceManager] load disk raw config(%v).", d) @@ -489,7 +520,12 @@ func (s *DataNode) startSpaceManager(cfg *config.Config) (err error) { wg.Add(1) go func(wg *sync.WaitGroup, path string, reservedSpace uint64) { defer wg.Done() - s.space.LoadDisk(path, reservedSpace, diskRdonlySpace, DefaultDiskMaxErr, diskEnableReadRepairExtentLimit) + if _, broken := brokenDisks[path]; !broken { + s.space.LoadDisk(path, reservedSpace, diskRdonlySpace, DefaultDiskMaxErr, diskEnableReadRepairExtentLimit) + return + } + log.LogWarnf("[startSpaceManager] load broken disk(%v)", path) + s.space.LoadBrokenDisk(path, reservedSpace, diskRdonlySpace, DefaultDiskMaxErr, diskEnableReadRepairExtentLimit) }(&wg, path, reservedSpace) } diff --git a/datanode/space_manager.go b/datanode/space_manager.go index 2e4147e86..00fb3be63 100644 --- a/datanode/space_manager.go +++ b/datanode/space_manager.go @@ -312,6 +312,22 @@ func (manager *SpaceManager) LoadDisk(path string, reservedSpace, diskRdonlySpac return } +func (manager *SpaceManager) LoadBrokenDisk(path string, reservedSpace, diskRdonlySpace uint64, maxErrCnt int, diskEnableReadRepairExtentLimit bool) (err error) { + var disk *Disk + + if diskRdonlySpace < reservedSpace { + diskRdonlySpace = reservedSpace + } + + log.LogDebugf("action[LoadBrokenDisk] load broken disk from path(%v).", path) + + if _, err = manager.GetDisk(path); err != nil { + disk = NewBrokenDisk(path, reservedSpace, diskRdonlySpace, maxErrCnt, manager, diskEnableReadRepairExtentLimit) + manager.putDisk(disk) + } + return +} + func (manager *SpaceManager) GetDisk(path string) (d *Disk, err error) { manager.diskMutex.RLock() defer manager.diskMutex.RUnlock() diff --git a/master/cluster_task.go b/master/cluster_task.go index 24ab8aa85..93a523671 100644 --- a/master/cluster_task.go +++ b/master/cluster_task.go @@ -1015,7 +1015,7 @@ func (c *Cluster) handleDataNodeHeartbeatResp(nodeAddr string, resp *proto.DataN dataNode.CpuUtil.Store(resp.CpuUtil) dataNode.SetIoUtils(resp.IoUtils) - dataNode.updateNodeMetric(resp) + dataNode.updateNodeMetric(c, resp) if err = c.t.putDataNode(dataNode); err != nil { log.LogErrorf("action[handleDataNodeHeartbeatResp] dataNode[%v],zone[%v],node set[%v], err[%v]", dataNode.Addr, dataNode.ZoneName, dataNode.NodeSetID, err) diff --git a/master/data_node.go b/master/data_node.go index af3906d11..0c3420243 100644 --- a/master/data_node.go +++ b/master/data_node.go @@ -16,6 +16,7 @@ package master import ( "fmt" + "sort" "sync" "sync/atomic" "time" @@ -147,7 +148,28 @@ func (dataNode *DataNode) getDisks(c *Cluster) (diskPaths []string) { return } -func (dataNode *DataNode) updateNodeMetric(resp *proto.DataNodeHeartbeatResponse) { +func (dataNode *DataNode) updateBadDisks(latest []string) (ok bool) { + sort.Slice(latest, func(i, j int) bool { + return latest[i] < latest[j] + }) + + curr := dataNode.BadDisks + dataNode.BadDisks = latest + if len(curr) != len(latest) { + ok = true + return + } + + for i := 0; i < len(curr); i++ { + if curr[i] != latest[i] { + ok = true + return + } + } + return +} + +func (dataNode *DataNode) updateNodeMetric(c *Cluster, resp *proto.DataNodeHeartbeatResponse) { dataNode.Lock() defer dataNode.Unlock() dataNode.DomainAddr = util.ParseIpAddrToDomainAddr(dataNode.Addr) @@ -164,7 +186,7 @@ func (dataNode *DataNode) updateNodeMetric(resp *proto.DataNodeHeartbeatResponse dataNode.TotalPartitionSize = resp.TotalPartitionSize dataNode.AllDisks = resp.AllDisks - dataNode.BadDisks = resp.BadDisks + updated := dataNode.updateBadDisks(resp.BadDisks) dataNode.BadDiskStats = resp.BadDiskStats dataNode.DiskStats = resp.DiskStats dataNode.BackupDataPartitions = resp.BackupDataPartitions @@ -177,6 +199,14 @@ func (dataNode *DataNode) updateNodeMetric(resp *proto.DataNodeHeartbeatResponse } dataNode.ReportTime = time.Now() dataNode.isActive = true + + if updated { + log.LogInfof("[updateNodeMetric] update data node(%v)", dataNode.Addr) + if err := c.syncUpdateDataNode(dataNode); err != nil { + log.LogErrorf("[updateNodeMetric] failed to update datanode(%v), err(%v)", dataNode.Addr, err) + } + } + log.LogDebugf("updateNodeMetric. datanode id %v addr %v total %v used %v avaliable %v", dataNode.ID, dataNode.Addr, dataNode.Total, dataNode.Used, dataNode.AvailableSpace) } diff --git a/master/metadata_fsm_op.go b/master/metadata_fsm_op.go index 7bbfb18c1..eeb05d1bb 100644 --- a/master/metadata_fsm_op.go +++ b/master/metadata_fsm_op.go @@ -421,6 +421,7 @@ type dataNodeValue struct { ToBeOffline bool DecommissionDiskList []string DecommissionDpTotal int + BadDisks []string } func newDataNodeValue(dataNode *DataNode) *dataNodeValue { @@ -439,6 +440,7 @@ func newDataNodeValue(dataNode *DataNode) *dataNodeValue { ToBeOffline: dataNode.ToBeOffline, DecommissionDiskList: dataNode.DecommissionDiskList, DecommissionDpTotal: dataNode.DecommissionDpTotal, + BadDisks: dataNode.BadDisks, } } @@ -1452,6 +1454,7 @@ func (c *Cluster) loadDataNodes() (err error) { dataNode.ToBeOffline = dnv.ToBeOffline dataNode.DecommissionDiskList = dnv.DecommissionDiskList dataNode.DecommissionDpTotal = dnv.DecommissionDpTotal + dataNode.BadDisks = dnv.BadDisks olddn, ok := c.dataNodes.Load(dataNode.Addr) if ok { if olddn.(*DataNode).ID <= dataNode.ID {