From 91c65e17aba16828b3d002c40930425f0d356e79 Mon Sep 17 00:00:00 2001 From: NaturalSelect Date: Fri, 12 Apr 2024 14:48:27 +0800 Subject: [PATCH] feat(data): speed data node restart Signed-off-by: NaturalSelect --- datanode/disk.go | 4 +- storage/extent.go | 53 ++++++ storage/extent_store.go | 337 ++++++++++++++++++++++++++++++---- util/fileutil/readdir.go | 27 +++ util/fileutil/readdir_test.go | 55 ++++++ util/fileutil/stat_test.go | 14 ++ 6 files changed, 451 insertions(+), 39 deletions(-) create mode 100644 util/fileutil/readdir.go create mode 100644 util/fileutil/readdir_test.go diff --git a/datanode/disk.go b/datanode/disk.go index 33bbee225..e54eb210d 100644 --- a/datanode/disk.go +++ b/datanode/disk.go @@ -49,7 +49,7 @@ const ( ExpiredPartitionExistTime = time.Hour * time.Duration(24*7) ) -const DefaultCurrentLoadDpLimit = 32 +const DefaultCurrentLoadDpLimit = 1 const ( DecommissionDiskMark = "decommissionDiskMark" @@ -626,7 +626,7 @@ func (d *Disk) RestorePartition(visitor PartitionVisitor, allowDelay bool) (err ) defer func() { if err == nil { - log.LogInfof("[RestorePartition] disk(%v) load dp(%v) using time(%v)", d.Path, dp.partitionID, time.Since(begin)) + log.LogInfof("[RestorePartition] disk(%v) load dp(%v) using time(%v) slow(%v)", d.Path, dp.partitionID, time.Since(begin), time.Since(begin) > 1*time.Second) } }() if dp, err = LoadDataPartition(path.Join(d.Path, filename), d, allowDelay); err != nil { diff --git a/storage/extent.go b/storage/extent.go index 0d1d10548..060e96c82 100644 --- a/storage/extent.go +++ b/storage/extent.go @@ -15,6 +15,7 @@ package storage import ( + "bytes" "encoding/binary" "fmt" "hash/crc32" @@ -71,6 +72,58 @@ func (ei *ExtentInfo) String() (m string) { return fmt.Sprintf("FileID(%v)_Size(%v)_IsDeleted(%v)_Souarce(%v)_MT(%d)_AT(%d)_CRC(%d)", ei.FileID, ei.Size, ei.IsDeleted, source, ei.ModifyTime, ei.AccessTime, ei.Crc) } +func (ei *ExtentInfo) MarshalBinaryWithBuffer(buff *bytes.Buffer) (err error) { + if err = binary.Write(buff, binary.BigEndian, ei.FileID); err != nil { + return + } + if err = binary.Write(buff, binary.BigEndian, ei.Size); err != nil { + return + } + if err = binary.Write(buff, binary.BigEndian, ei.IsDeleted); err != nil { + return + } + if err = binary.Write(buff, binary.BigEndian, ei.ModifyTime); err != nil { + return + } + if err = binary.Write(buff, binary.BigEndian, ei.AccessTime); err != nil { + return + } + return +} + +func (ei *ExtentInfo) MarshalBinary() (v []byte, err error) { + buff := bytes.NewBuffer([]byte{}) + if err = ei.MarshalBinaryWithBuffer(buff); err != nil { + return + } + v = buff.Bytes() + return +} + +func (ei *ExtentInfo) UnmarshalBinaryWithBuffer(buff *bytes.Buffer) (err error) { + if err = binary.Read(buff, binary.BigEndian, &ei.FileID); err != nil { + return + } + if err = binary.Read(buff, binary.BigEndian, &ei.Size); err != nil { + return + } + if err = binary.Read(buff, binary.BigEndian, &ei.IsDeleted); err != nil { + return + } + if err = binary.Read(buff, binary.BigEndian, &ei.ModifyTime); err != nil { + return + } + if err = binary.Read(buff, binary.BigEndian, &ei.AccessTime); err != nil { + return + } + return +} + +func (ei *ExtentInfo) UnmarshalBinary(v []byte) (err error) { + err = ei.UnmarshalBinaryWithBuffer(bytes.NewBuffer(v)) + return +} + // SortedExtentInfos defines an array sorted by AccessTime type SortedExtentInfos []*ExtentInfo diff --git a/storage/extent_store.go b/storage/extent_store.go index a3ea1d3e4..2ca7b439f 100644 --- a/storage/extent_store.go +++ b/storage/extent_store.go @@ -39,6 +39,7 @@ import ( "github.com/cubefs/cubefs/proto" "github.com/cubefs/cubefs/util" "github.com/cubefs/cubefs/util/errors" + "github.com/cubefs/cubefs/util/fileutil" "github.com/cubefs/cubefs/util/log" "github.com/cubefs/cubefs/util/strutil" ) @@ -62,6 +63,8 @@ const ( NormalExtentDeleteRetainTime = 3600 * 4 CacheFlushInterval = 5 * time.Second + ExtentReadDirHint = "READDIR_HINT" + ExtentReadDirHintTemp = "READDIR_HINT.tmp" StaleExtStoreBackupSuffix = ".old" StaleExtStoreTimeFormat = "20060102150405.000000000" @@ -109,6 +112,39 @@ var ( } ) +var delayMark *ExtentInfo = &ExtentInfo{} + +type ExtentInfoOnDemand struct { + ei atomic.Value +} + +func (eiod *ExtentInfoOnDemand) IsLoaded() (ok bool) { + return eiod.ei.Load() != delayMark +} + +func (eiod *ExtentInfoOnDemand) Load(ei *ExtentInfo) { + eiod.ei.CompareAndSwap(delayMark, ei) +} + +func (eiod *ExtentInfoOnDemand) Get() (ei *ExtentInfo) { + v := eiod.ei.Load().(*ExtentInfo) + if v != delayMark { + ei = v + } + return +} + +func NewExtentInfoOnDemand() (eiod *ExtentInfoOnDemand) { + eiod = NewExtentInfoOnDemandWithInfo(delayMark) + return +} + +func NewExtentInfoOnDemandWithInfo(ei *ExtentInfo) (eiod *ExtentInfoOnDemand) { + eiod = &ExtentInfoOnDemand{} + eiod.ei.Store(ei) + return +} + // ExtentStore defines fields used in the storage engine. // Packets smaller than 128K are stored in the "tinyExtent", a place to persist the small files. // packets larger than or equal to 128K are stored in the normal "extent", a place to persist large files. @@ -242,7 +278,7 @@ func NewExtentStore(dataDir string, partitionID uint64, storeSize, dpType int, i s.extentInfoMap = make(map[uint64]*ExtentInfo) s.extentLockMap = make(map[uint64]proto.GcFlag, 0) s.cache = NewExtentCache(100) - if err = s.initBaseFileID(); err != nil { + if err = s.initBaseFileID(allowDelay, 200*time.Millisecond); err != nil { err = fmt.Errorf("init base field ID: %v", err) return } @@ -254,7 +290,10 @@ func NewExtentStore(dataDir string, partitionID uint64, storeSize, dpType int, i return } s.stopC = make(chan interface{}) - go s.startFlushCache() + go func() { + time.Sleep(15 * time.Minute) + s.startFlushCache() + }() return } @@ -358,56 +397,276 @@ func (s *ExtentStore) Create(extentID uint64) (err error) { return } -func (s *ExtentStore) initBaseFileID() error { +func (s *ExtentStore) GetExtentInfoFromDisk(id uint64) (ei *ExtentInfo, err error) { + retry := 0 + const maxRetry = 3 + + for retry < maxRetry { + var stat fs.FileInfo + name := path.Join(s.dataPath, fmt.Sprint(id)) + stat, err = os.Stat(name) + if err != nil { + retry++ + continue + } + + ino := stat.Sys().(*syscall.Stat_t) + ei = &ExtentInfo{ + FileID: id, + Size: uint64(stat.Size()), + Crc: 0, + IsDeleted: false, + AccessTime: time.Unix(int64(ino.Atim.Sec), int64(ino.Atim.Nsec)).Unix(), + ModifyTime: stat.ModTime().Unix(), + Source: "", + } + if IsTinyExtent(id) { + watermark := ei.Size + if watermark%PageSize != 0 { + watermark = watermark + (PageSize - watermark%PageSize) + } + ei.Size = watermark + } + return + } + return +} + +func (s *ExtentStore) GetExtentInfoFromMap(id uint64) (ei *ExtentInfo, ok bool) { + s.eiMutex.RLock() + defer s.eiMutex.RUnlock() + v, ok := s.extentInfoMap[id] + if !ok { + return + } + ei = v.Get() + return +} + +func (s *ExtentStore) GetExtentInfo(id uint64) (ei *ExtentInfo, ok bool, err error) { + s.eiMutex.RLock() + defer s.eiMutex.RUnlock() + + v, ok := s.extentInfoMap[id] + if !ok { + return + } + + if !v.IsLoaded() { + ei, err = s.GetExtentInfoFromDisk(id) + if err != nil { + log.LogErrorf("[GetExtentInfo] failed to load extent(%v) info, err(%v)", id, err) + return + } + v.Load(ei) + ok = true + } + ei = v.Get() + return +} + +func (s *ExtentStore) SetExtentInfo(id uint64, ei *ExtentInfo) { + s.eiMutex.Lock() + defer s.eiMutex.Unlock() + v := NewExtentInfoOnDemand() + v.Load(ei) + s.extentInfoMap[id] = v +} + +func (s *ExtentStore) SetExtentInfoDelay(id uint64) { + s.eiMutex.Lock() + defer s.eiMutex.Unlock() + s.extentInfoMap[id] = NewExtentInfoOnDemand() +} + +func (s *ExtentStore) RangeExtentInfo(iter func(id uint64, ei *ExtentInfo) (ok bool, err error)) (err error) { + s.eiMutex.RLock() + defer s.eiMutex.RUnlock() + + var ok bool + for id, v := range s.extentInfoMap { + if !v.IsLoaded() { + var ei *ExtentInfo + ei, err = s.GetExtentInfoFromDisk(id) + if err != nil { + log.LogErrorf("[RangeExtentInfo] failed to load extent(%v) info, err(%v)", id, err) + return + } + v.Load(ei) + } + + ei := v.Get() + ok, err = iter(id, ei) + if err != nil || !ok { + return + } + } + return +} + +func (s *ExtentStore) DeleteExtentInfo(id uint64) { + s.eiMutex.Lock() + defer s.eiMutex.Unlock() + delete(s.extentInfoMap, id) +} + +func (s *ExtentStore) GetExtentInfoCount() (count int) { + s.eiMutex.RLock() + defer s.eiMutex.RUnlock() + count = len(s.extentInfoMap) + return +} + +func (s *ExtentStore) writeReadDirHint() (err error) { + hintTempPath := path.Join(s.dataPath, ExtentReadDirHintTemp) + hintPath := path.Join(s.dataPath, ExtentReadDirHint) + buff := bytes.NewBuffer([]byte{}) + err = s.RangeExtentInfo(func(id uint64, ei *ExtentInfo) (ok bool, err error) { + err = ei.MarshalBinaryWithBuffer(buff) + if err != nil { + return + } + return true, nil + }) + if err != nil { + log.LogErrorf("[writeReadDirHint] store(%v) failed to marshal hint, err(%v)", s.dataPath, err) + return + } + if err = os.WriteFile(hintTempPath, buff.Bytes(), 0666); err != nil { + log.LogErrorf("[writeReadDirHint] store(%v) failed to write readdir hint, err(%v)", s.dataPath, err) + return + } + err = os.Rename(hintTempPath, hintPath) + if err != nil { + log.LogErrorf("[writeReadDirHint] store(%v) failed to rename readdir hint, err(%v)", s.dataPath, err) + return + } + return +} + +func (s *ExtentStore) readReadDirHint() (extMap map[uint64]*ExtentInfoOnDemand, err error) { + var data []byte + begin := time.Now() + defer func() { + size := 0 + cnt := 0 + if data != nil { + size = len(data) + } + if extMap != nil { + cnt = len(extMap) + } + slow := time.Since(begin) > 1*time.Second + log.LogInfof("[readReadDirHint] store(%v) read hint file using time(%v), read size(%v), cnt(%v), slow(%v)", s.dataPath, time.Since(begin), size, cnt, slow) + }() + + hintPath := path.Join(s.dataPath, ExtentReadDirHint) + data, err = os.ReadFile(hintPath) + if err != nil { + if os.IsNotExist(err) { + err = nil + return + } + log.LogErrorf("[readReadDirHint] store(%v) failed to read hint file, err(%v)", s.dataPath, err) + return + } + extMap = make(map[uint64]*ExtentInfoOnDemand) + buff := bytes.NewBuffer(data) + for buff.Len() != 0 { + ei := &ExtentInfo{} + err = ei.UnmarshalBinaryWithBuffer(buff) + if err != nil { + log.LogErrorf("[readReadDirHint] store(%v) failed to unmarshal hint, err(%v)", s.dataPath, err) + return + } + eiod := NewExtentInfoOnDemandWithInfo(ei) + extMap[ei.FileID] = eiod + } + return +} + +func (s *ExtentStore) removeReadDirHint() (err error) { + hintPath := path.Join(s.dataPath, ExtentReadDirHint) + err = os.Remove(hintPath) + if err != nil { + if os.IsNotExist(err) { + err = nil + return + } + log.LogErrorf("[removeReadDirHint] store(%v) failed to remove read dir hint, err(%v)", s.dataPath, err) + return + } + return +} + +func (s *ExtentStore) initBaseFileID(allowDelay bool, loadTimeout time.Duration) error { var extNum int begin := time.Now() defer func() { - log.LogInfof("[initBaseFileID] init base file id using time(%v), count(%v)", time.Since(begin), extNum) + log.LogInfof("[initBaseFileID] store(%v) init base file id using time(%v), count(%v)", s.dataPath, time.Since(begin), extNum) }() var ( baseFileID uint64 ) baseFileID, _ = s.GetPersistenceBaseExtentID() - files, err := os.ReadDir(s.dataPath) + log.LogInfof("[initBaseFileID] store(%v) init base file to persistence base extent id using time(%v)", s.dataPath, time.Since(begin)) + + // NOTE: try to read hint + var err error + var extMap map[uint64]*ExtentInfoOnDemand + extMap, err = s.readReadDirHint() if err != nil { + log.LogErrorf("[initBaseFileID] store(%v) failed to read hint, err(%v)", s.dataPath, err) + err = nil + } + // NOTE: remove hint + if err = s.removeReadDirHint(); err != nil { + log.LogErrorf("[initBaseFileID] store(%v) failed to remove hint, err(%v)", s.dataPath, err) return err } - var ( - e *Extent - ei *ExtentInfo - loadErr error - ) - extentIds := make([]uint64, 0, len(files)) - for _, f := range files { - if extentID, isExtent := s.ExtentID(f.Name()); isExtent { - extentIds = append(extentIds, extentID) + if len(extMap) != 0 { + log.LogInfof("[initBaseFileID] store(%v) init base file to read hint using time(%v)", s.dataPath, time.Since(begin)) + // NOTE: fast path + for id := range extMap { + if !IsTinyExtent(id) && id > baseFileID { + baseFileID = id + } + } + s.extentInfoMap = extMap + } else { + // NOTE: slow path + files, err := fileutil.ReadDir(s.dataPath) + if err != nil { + return err + } + log.LogInfof("[initBaseFileID] store(%v) init base file to read dir using time(%v)", s.dataPath, time.Since(begin)) + + for _, f := range files { + extentID, isExtent := s.ExtentID(f) + if !isExtent { + continue + } + extNum++ + s.SetExtentInfoDelay(extentID) + // NOTE: if not timeout, load extent info + if time.Since(begin) < loadTimeout || !allowDelay { + _, _, err = s.GetExtentInfo(extentID) + if err != nil { + log.LogErrorf("[initBaseFileID] store(%v) failed to load extent(%v), err(%v)", s.dataPath, extentID, err) + return err + } + } else { + log.LogInfof("[initBaseFileID] store(%v) load using time(%v) too long, switch to on demand mode, loaded count(%v)", s.dataPath, time.Since(begin), extNum) + } + + if !IsTinyExtent(extentID) && extentID > baseFileID { + baseFileID = extentID + } } } - sort.Slice(extentIds, func(i, j int) bool { - return extentIds[i] < extentIds[j] - }) - - for _, extentID := range extentIds { - if e, loadErr = s.extent(extentID); loadErr != nil { - log.LogError("[initBaseFileID] load extent error", loadErr) - continue - } - - ei = &ExtentInfo{FileID: extentID} - ei.UpdateExtentInfo(e, 0) - atomic.StoreInt64(&ei.AccessTime, e.accessTime) - - s.eiMutex.Lock() - s.extentInfoMap[extentID] = ei - s.eiMutex.Unlock() - - e.Close() - if !IsTinyExtent(extentID) && extentID > baseFileID { - baseFileID = extentID - } - } + log.LogInfof("[initBaseFileID] store(%v) init base file to load loop using time(%v)", s.dataPath, time.Since(begin)) if baseFileID < MinExtentID { baseFileID = MinExtentID } @@ -727,6 +986,10 @@ func (s *ExtentStore) Close() { } } s.closed = true + + if err := s.writeReadDirHint(); err != nil { + log.LogErrorf("[Close] store(%v) failed to write extent hint, err(%v)", s.dataPath, err) + } } // Watermark returns the extent info of the given extent on the record. diff --git a/util/fileutil/readdir.go b/util/fileutil/readdir.go new file mode 100644 index 000000000..0093f2306 --- /dev/null +++ b/util/fileutil/readdir.go @@ -0,0 +1,27 @@ +// Copyright 2024 The CubeFS Authors. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or +// implied. See the License for the specific language governing +// permissions and limitations under the License. + +package fileutil + +import "os" + +func ReadDir(name string) (dentries []string, err error) { + dir, err := os.Open(name) + if err != nil { + return + } + defer dir.Close() + dentries, err = dir.Readdirnames(0) + return +} diff --git a/util/fileutil/readdir_test.go b/util/fileutil/readdir_test.go new file mode 100644 index 000000000..a0dab8075 --- /dev/null +++ b/util/fileutil/readdir_test.go @@ -0,0 +1,55 @@ +// Copyright 2024 The CubeFS Authors. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or +// implied. See the License for the specific language governing +// permissions and limitations under the License. + +package fileutil_test + +import ( + "os" + "path" + "testing" + + "github.com/cubefs/cubefs/util/fileutil" + "github.com/stretchr/testify/require" +) + +func TestReaddir(t *testing.T) { + tmpDir, err := os.MkdirTemp("", "") + require.NoError(t, err) + + file1 := "f1" + file2 := "f2" + + err = os.WriteFile(path.Join(tmpDir, file1), []byte(""), 0644) + require.NoError(t, err) + + os.WriteFile(path.Join(tmpDir, file2), []byte(""), 0644) + require.NoError(t, err) + + names, err := fileutil.ReadDir(tmpDir) + require.NoError(t, err) + + fileMap := make(map[string]interface{}) + + for _, v := range names { + fileMap[v] = 1 + } + + require.EqualValues(t, 2, len(fileMap)) + + _, file1Exist := fileMap[file1] + _, file2Exist := fileMap[file2] + + require.True(t, file1Exist) + require.True(t, file2Exist) +} diff --git a/util/fileutil/stat_test.go b/util/fileutil/stat_test.go index 1d8b0b331..82da461a1 100644 --- a/util/fileutil/stat_test.go +++ b/util/fileutil/stat_test.go @@ -1,3 +1,17 @@ +// Copyright 2024 The CubeFS Authors. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or +// implied. See the License for the specific language governing +// permissions and limitations under the License. + package fileutil_test import (