feat(data): speed data node restart

Signed-off-by: NaturalSelect <huangzhibin1@oppo.com>
This commit is contained in:
NaturalSelect 2024-04-12 14:48:27 +08:00 committed by longerfly
parent 1e9dec3c20
commit 91c65e17ab
6 changed files with 451 additions and 39 deletions

View File

@ -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 {

View File

@ -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

View File

@ -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.

27
util/fileutil/readdir.go Normal file
View File

@ -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
}

View File

@ -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)
}

View File

@ -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 (