feat(data): limit io current on a disk when load/stop dp

Signed-off-by: S9054862 <huangzhibin1@oppo.com>
This commit is contained in:
S9054862 2024-03-19 14:16:37 +08:00 committed by longerfly
parent 0c6f5314da
commit 00ecb62b94
8 changed files with 218 additions and 26 deletions

View File

@ -49,6 +49,8 @@ const (
ExpiredPartitionExistTime = time.Hour * time.Duration(24*7)
)
const DefaultCurrentLoadDpLimit = 8
const (
DecommissionDiskMark = "decommissionDiskMark"
)
@ -560,6 +562,11 @@ func (d *Disk) isExpiredPartitionDir(filename string) (isExpiredPartitionDir boo
return
}
type dpLoadInfo struct {
Id uint64
FileName string
}
// RestorePartition reads the files stored on the local disk and restores the data partitions.
func (d *Disk) RestorePartition(visitor PartitionVisitor) (err error) {
convert := func(node *proto.DataNodeInfo) *DataNodeInfo {
@ -597,7 +604,54 @@ func (d *Disk) RestorePartition(visitor PartitionVisitor) (err error) {
var (
wg sync.WaitGroup
toDeleteExpiredPartitionNames = make([]string, 0)
dpNum int
)
begin := time.Now()
defer func() {
msg := fmt.Sprintf("[RestorePartition] disk(%v) load all dp(%v) using time(%v)", d.Path, dpNum, time.Since(begin))
syslog.Print(msg)
log.LogInfo(msg)
}()
loadCh := make(chan dpLoadInfo, d.space.currentLoadDpCount)
for i := 0; i < d.space.currentLoadDpCount; i++ {
wg.Add(1)
go func() {
defer wg.Done()
loader := func(partitionID uint64, filename string) {
begin := time.Now()
var (
dp *DataPartition
err error
)
defer func() {
if err == nil {
log.LogInfof("[RestorePartition] disk(%v) load dp(%v) using time(%v)", d.Path, dp.partitionID, time.Since(begin))
}
}()
if dp, err = LoadDataPartition(path.Join(d.Path, filename), d); err != nil {
mesg := fmt.Sprintf("action[RestorePartition] new partition(%v) err(%v) ",
partitionID, err.Error())
log.LogError(mesg)
exporter.Warning(mesg)
syslog.Println(mesg)
return
}
if visitor != nil {
visitor(dp)
}
}
for {
dp, ok := <-loadCh
if !ok {
return
}
loader(dp.Id, dp.FileName)
}
}()
}
for _, fileInfo := range fileInfoList {
filename := fileInfo.Name()
if !d.isPartitionDir(filename) {
@ -626,31 +680,13 @@ func (d *Disk) RestorePartition(visitor PartitionVisitor) (err error) {
toDeleteExpiredPartitionNames = append(toDeleteExpiredPartitionNames, newName)
continue
}
wg.Add(1)
go func(partitionID uint64, filename string) {
var (
dp *DataPartition
err error
)
defer wg.Done()
if dp, err = LoadDataPartition(path.Join(d.Path, filename), d); err != nil {
mesg := fmt.Sprintf("action[RestorePartition] new partition(%v) err(%v) ",
partitionID, err.Error())
log.LogError(mesg)
if IsDiskErr(err.Error()) {
d.triggerDiskError(ReadFlag, partitionID)
}
// exporter.Warning(mesg)
// syslog.Println(mesg)
return
}
if visitor != nil {
visitor(dp)
}
}(partitionID, filename)
dpNum++
loadCh <- dpLoadInfo{
Id: partitionID,
FileName: filename,
}
}
close(loadCh)
if len(toDeleteExpiredPartitionNames) > 0 {
log.LogInfof("action[RestorePartition] expiredPartitions %v, disk %v", toDeleteExpiredPartitionNames, d.Path)

View File

@ -285,7 +285,15 @@ func LoadDataPartition(partitionDir string, disk *Disk) (dp *DataPartition, err
dp.lastTruncateID = meta.LastTruncateID
if meta.DataPartitionCreateType == proto.NormalCreateDataPartition {
err = dp.StartRaft(true)
func() {
begin := time.Now()
defer func() {
log.LogInfof("[LoadDataPartition] load dp(%v) flush extent using time(%v)", dp.partitionID, time.Since(begin))
}()
dp.extentStore.Flush()
}()
} else {
log.LogInfof("[LoadDataPartition] dp(%v) skip disk limit, need repair", dp.partitionID)
// init leaderSize to partitionSize
dp.leaderSize = dp.partitionSize
dp.partitionStatus = proto.Recovering
@ -320,6 +328,11 @@ func LoadDataPartition(partitionDir string, disk *Disk) (dp *DataPartition, err
func newDataPartition(dpCfg *dataPartitionCfg, disk *Disk, isCreate bool) (dp *DataPartition, err error) {
partitionID := dpCfg.PartitionID
begin := time.Now()
defer func() {
log.LogInfof("[newDataPartition] load dp(%v) new data partition using time(%v)", partitionID, time.Since(begin))
}()
var dataPath string
if proto.IsNormalDp(dpCfg.PartitionType) {
@ -611,6 +624,11 @@ func (dp *DataPartition) SnapShot() (files []*proto.File) {
// Stop close the store and the raft store.
func (dp *DataPartition) Stop() {
begin := time.Now()
defer func() {
msg := fmt.Sprintf("[Stop] stop dp(%v) using time(%v)", dp.partitionID, time.Since(begin))
log.LogInfo(msg)
}()
dp.stopOnce.Do(func() {
if dp.stopC != nil {
close(dp.stopC)

View File

@ -18,6 +18,8 @@ import (
"encoding/binary"
"encoding/json"
"fmt"
"io/ioutil"
syslog "log"
"net"
"os"
"path"
@ -77,6 +79,11 @@ func (dp *DataPartition) raftPort() (heartbeat, replica int, err error) {
// StartRaft start raft instance when data partition start or restore.
func (dp *DataPartition) StartRaft(isLoad bool) (err error) {
begin := time.Now()
defer func() {
log.LogInfof("[StartRaft] load dp(%v) start raft using time(%v)", dp.partitionID, time.Since(begin))
}()
// cache or preload partition not support raft and repair.
if !dp.isNormalType() {
return nil
@ -138,6 +145,10 @@ func (dp *DataPartition) raftStopped() bool {
}
func (dp *DataPartition) stopRaft() {
begin := time.Now()
defer func() {
log.LogInfof("[stopRaft] dp(%v) stop raft using time(%v)", dp.partitionID, time.Since(begin))
}()
if atomic.CompareAndSwapInt32(&dp.raftStatus, RaftStatusRunning, RaftStatusStopped) {
// cache or preload partition not support raft and repair.
if !dp.isNormalType() {
@ -473,6 +484,10 @@ func (dp *DataPartition) storeAppliedID(applyIndex uint64) (err error) {
// LoadAppliedID loads the applied IDs to the memory.
func (dp *DataPartition) LoadAppliedID() (err error) {
begin := time.Now()
defer func() {
log.LogInfof("[LoadAppliedID] load dp(%v) load applied id using time(%v)", dp.partitionID, time.Since(begin))
}()
filename := path.Join(dp.Path(), ApplyIndexFile)
if _, err = os.Stat(filename); err != nil {
return
@ -580,6 +595,12 @@ func (s *DataNode) startRaftServer(cfg *config.Config) (err error) {
func (s *DataNode) stopRaftServer() {
if s.raftStore != nil {
begin := time.Now()
defer func() {
msg := fmt.Sprintf("[stopRaftServer] stop raft server using time(%v)", time.Since(begin))
log.LogInfo(msg)
syslog.Print(msg)
}()
s.raftStore.Stop()
}
}

View File

@ -29,6 +29,11 @@ import (
"strings"
"sync"
"sync/atomic"
"time"
"errors"
syslog "log"
"os"
"syscall"
"time"
@ -118,6 +123,10 @@ const (
ConfigDiskWriteIops = "diskWriteIops" // int
ConfigDiskWriteFlow = "diskWriteFlow" // int
// load/stop dp limit
ConfigDiskCurrentLoadDpLimit = "diskCurrentLoadDpLimit"
ConfigDiskCurrentStopDpLimit = "diskCurrentStopDpLimit"
ConfigServiceIDKey = "serviceIDKey"
// disk status becomes unavailable if disk error partition count reaches this value
@ -280,6 +289,12 @@ func doStart(server common.Server, cfg *config.Config) (err error) {
}
func doShutdown(server common.Server) {
begin := time.Now()
defer func() {
msg := fmt.Sprintf("[doShutdown] stop datanode using time(%v)", time.Since(begin))
log.LogInfo(msg)
syslog.Print(msg)
}()
s, ok := server.(*DataNode)
if !ok {
return
@ -396,6 +411,11 @@ func (s *DataNode) startSpaceManager(cfg *config.Config) (err error) {
s.space.SetClusterID(s.clusterID)
s.initQosLimit(cfg)
loadLimit := cfg.GetInt(ConfigDiskCurrentLoadDpLimit)
stopLimit := cfg.GetInt(ConfigDiskCurrentStopDpLimit)
s.space.SetCurrentLoadDpLimit(loadLimit)
s.space.SetCurrentStopDpLimit(stopLimit)
diskRdonlySpace := uint64(cfg.GetInt64(CfgDiskRdonlySpace))
if diskRdonlySpace < DefaultDiskRetainMin {
diskRdonlySpace = DefaultDiskRetainMin
@ -690,7 +710,12 @@ func (s *DataNode) startTCPService() (err error) {
func (s *DataNode) stopTCPService() (err error) {
if s.tcpListener != nil {
begin := time.Now()
defer func() {
msg := fmt.Sprintf("[stopTCPService] stop tcp service using time(%v)", time.Since(begin))
log.LogInfo(msg)
syslog.Print(msg)
}()
s.tcpListener.Close()
log.LogDebugf("action[stopTCPService] stop tcp service.")
}
@ -741,6 +766,12 @@ func (s *DataNode) startSmuxService(cfg *config.Config) (err error) {
func (s *DataNode) stopSmuxService() (err error) {
if s.smuxListener != nil {
begin := time.Now()
defer func() {
msg := fmt.Sprintf("[stopSmuxService] stop smux service uing time(%v)", time.Since(begin))
syslog.Print(msg)
log.LogInfo(msg)
}()
s.smuxListener.Close()
log.LogDebugf("action[stopSmuxService] stop smux service.")
}
@ -866,6 +897,12 @@ func (s *DataNode) initConnPool() {
func (s *DataNode) closeSmuxConnPool() {
if s.smuxConnPool != nil {
begin := time.Now()
defer func() {
msg := fmt.Sprintf("[closeSmuxConnPool] close smux conn pool using time(%v)", time.Since(begin))
log.LogInfo(msg)
syslog.Print(msg)
}()
s.smuxConnPool.Close()
log.LogDebugf("action[stopSmuxService] stop smux conn pool")
}

View File

@ -24,6 +24,11 @@ import (
"sync/atomic"
"time"
"math"
"os"
syslog "log"
"github.com/cubefs/cubefs/proto"
"github.com/cubefs/cubefs/raftstore"
"github.com/cubefs/cubefs/util"
@ -33,6 +38,8 @@ import (
"github.com/shirou/gopsutil/disk"
)
const DefaultStopDpLimit = DefaultCurrentLoadDpLimit
// SpaceManager manages the disk space.
type SpaceManager struct {
clusterID string
@ -49,6 +56,8 @@ type SpaceManager struct {
dataNode *DataNode
createPartitionMutex sync.RWMutex
rand *rand.Rand
currentLoadDpCount int
currentStopDpCount int
diskUtils map[string]*atomicutil.Float64
samplerDone chan struct{}
allDisksLoaded bool
@ -66,13 +75,33 @@ func NewSpaceManager(dataNode *DataNode) *SpaceManager {
space.stopC = make(chan bool)
space.dataNode = dataNode
space.rand = rand.New(rand.NewSource(time.Now().Unix()))
space.currentLoadDpCount = DefaultCurrentLoadDpLimit
space.currentStopDpCount = DefaultStopDpLimit
space.diskUtils = make(map[string]*atomicutil.Float64)
go space.statUpdateScheduler()
return space
}
func (manager *SpaceManager) SetCurrentLoadDpLimit(limit int) {
if limit != 0 {
manager.currentLoadDpCount = limit
}
}
func (manager *SpaceManager) SetCurrentStopDpLimit(limit int) {
if limit != 0 {
manager.currentStopDpCount = limit
}
}
func (manager *SpaceManager) Stop() {
begin := time.Now()
defer func() {
msg := fmt.Sprintf("[Stop] stop space manager using time(%v)", time.Since(begin))
log.LogInfo(msg)
syslog.Print(msg)
}()
defer func() {
recover()
}()
@ -91,6 +120,18 @@ func (manager *SpaceManager) Stop() {
partition.stopRaft()
}
limitCh := make(map[string]chan interface{})
for _, partition := range manager.partitions {
_, ok := limitCh[partition.Disk().Path]
if !ok {
ch := make(chan interface{}, manager.currentStopDpCount)
defer func() {
close(ch)
}()
limitCh[partition.Disk().Path] = ch
}
}
go func(c chan<- *DataPartition) {
defer wg.Done()
for _, partition := range manager.partitions {
@ -108,7 +149,14 @@ func (manager *SpaceManager) Stop() {
if partition = <-c; partition == nil {
return
}
partition.Stop()
func() {
limit := limitCh[partition.Disk().Path]
defer func() {
<-limit
}()
limit <- 1
partition.Stop()
}()
}
}(partitionC)
}
@ -388,6 +436,10 @@ func (manager *SpaceManager) Partition(partitionID uint64) (dp *DataPartition) {
}
func (manager *SpaceManager) AttachPartition(dp *DataPartition) {
begin := time.Now()
defer func() {
log.LogInfof("[AttachPartition] load dp(%v) attach using time(%v)", dp.partitionID, time.Now().Sub(begin))
}()
manager.partitionMutex.Lock()
defer manager.partitionMutex.Unlock()
manager.partitions[dp.partitionID] = dp

View File

@ -22,6 +22,9 @@
| diskWriteIocc | int | 限制单盘并发写操作,小于等于0表示不限制 | 否 |
| diskWriteFlow | int | 限制单盘写流量,小于等于0表示不限制 | 否 |
| disks | string slice | 格式:`磁盘挂载路径:预留空间` ,预留空间配置范围`[20G,50G]` | 是 |
| diskCurrentLoadDpLimit | int | 一个磁盘上并发加载的data partition的最大数量 | No |
| diskCurrentStopDpLimit | int | 一个磁盘上并发停止的data partition的最大数量 | No |
| enableLogPanicHook | bool | (实验性) Hook `panic` 函数以便在执行`panic`之前使日志落盘 | No | false |
## 配置示例

View File

@ -22,6 +22,9 @@
| diskWriteIocc | int | Limit write concurrency io frequency per disk. No limit if less than or equal to 0 | No |
| diskWriteFlow | int | Limit write io flow per disk. No limit if less than or equal to 0 | No |
| disks | string slice | Format: `disk mount path:reserved space`, reserved space configuration range `[20G,50G]` | Yes |
| diskCurrentLoadDpLimit | int | The max count of data partition on a disk that current load | No |
| diskCurrentStopDpLimit | int | The max count of data partition on a disk that current stop | No |
| enableLogPanicHook | bool | (Experimental) Hook `panic` function to flush log before executing `panic` | No | false |
## Configuration Example

View File

@ -148,6 +148,10 @@ func MkdirAll(name string) (err error) {
}
func NewExtentStore(dataDir string, partitionID uint64, storeSize, dpType int, isCreate bool) (s *ExtentStore, err error) {
begin := time.Now()
defer func() {
log.LogInfof("[NewExtentStore] load dp(%v) new extent store using time(%v)", partitionID, time.Since(begin))
}()
s = new(ExtentStore)
s.dataPath = dataDir
s.partitionType = dpType
@ -199,6 +203,7 @@ func NewExtentStore(dataDir string, partitionID uint64, storeSize, dpType int, i
needWriteEmpty := DeleteTinyRecordSize - (stat.Size() % DeleteTinyRecordSize)
data := make([]byte, needWriteEmpty)
s.tinyExtentDeleteFp.Write(data)
log.LogInfof("[NewExtentStore] load dp(%v) write zero buffer", partitionID)
}
log.LogDebugf("NewExtentStore.partitionID [%v] dataPath %v verifyExtentFp init", partitionID, s.dataPath)
@ -634,8 +639,25 @@ func (s *ExtentStore) IsDeletedNormalExtent(extentID uint64) (ok bool) {
return
}
func (s *ExtentStore) Flush() {
begin := time.Now()
defer func() {
log.LogInfof("[Flush] flush extent store using time(%v)", time.Since(begin))
}()
s.mutex.Lock()
defer s.mutex.Unlock()
if s.closed {
return
}
s.cache.Flush()
}
// Close closes the extent store.
func (s *ExtentStore) Close() {
begin := time.Now()
defer func() {
log.LogInfof("[Close] close extent store using time(%v)", time.Since(begin))
}()
s.mutex.Lock()
defer s.mutex.Unlock()
if s.closed {