fix(shardnode): optimize shard do checkpoint

with #22357426

Signed-off-by: xiejian <xiejian3@oppo.com>
This commit is contained in:
xiejian 2025-02-11 20:00:11 +08:00 committed by slasher
parent ad2ace90d8
commit 0269a61764
7 changed files with 110 additions and 53 deletions

View File

@ -30,6 +30,11 @@ const (
LevelStyle = CompactionStyle("level")
UniversalStyle = CompactionStyle("universal")
ReadTierAll = rdb.ReadTier(0)
ReadTierBlockCache = rdb.ReadTier(1)
ReadTierPersisted = rdb.ReadTier(2)
ReadTierMemtable = rdb.ReadTier(3)
defaultReadConcurrency = 10
defaultReadQueueLen = 10
defaultWriteConcurrency = 4
@ -90,6 +95,7 @@ type (
}
ReadOption interface {
SetSnapShot(snap Snapshot)
SetReadTier(tier rdb.ReadTier)
Close()
}
ReadOptFunc func(opts *readOpts)

View File

@ -234,7 +234,7 @@ func (s *rocksdb) writeLoop(ch chan *writeTask) {
case batchEvent:
tasks[i] = nil
b := task.batch
err := s.write(task.ctx, b, nil)
err := s.write(task.ctx, b, task.wo.opt)
task.err <- err
s.reduceWriteReqCnt(1)
default:
@ -347,6 +347,10 @@ func (ro *readOption) SetSnapShot(snap Snapshot) {
ro.opt.SetSnapshot(ro.snap)
}
func (ro *readOption) SetReadTier(tier rdb.ReadTier) {
ro.opt.SetReadTier(tier)
}
func (ro *readOption) Close() {
ro.opt.Destroy()
}
@ -441,6 +445,8 @@ func (lr *listReader) ReadNext() (key KeyGetter, val ValueGetter, err error) {
vg := &valueGetter{value: lr.iterator.Value()}
lr.isFirst = false
if lr.filterKey(kg) {
kg.Close()
vg.Close()
return lr.ReadNext()
}
return kg, vg, nil
@ -483,6 +489,8 @@ func (lr *listReader) ReadPrev() (key KeyGetter, val ValueGetter, err error) {
vg := &valueGetter{value: lr.iterator.Value()}
lr.isFirst = false
if lr.filterKey(kg) {
kg.Close()
vg.Close()
return lr.ReadPrev()
}
return kg, vg, nil
@ -691,13 +699,12 @@ func (s *rocksdb) CheckColumns(col CF) bool {
func (s *rocksdb) Get(ctx context.Context, col CF, key []byte, opts ...ReadOptFunc) (value ValueGetter, err error) {
ro := &readOpts{}
ro.applyOptions(opts)
if ro.withNoMerge {
if ro.opt != nil || ro.withNoMerge {
return s.get(ctx, col, key, ro.opt)
}
task := s.newReadTask(ctx)
task.typ = cfGet
task.ro = ro
task.cf = []CF{col}
task.key = [][]byte{key}
@ -719,13 +726,12 @@ func (s *rocksdb) Get(ctx context.Context, col CF, key []byte, opts ...ReadOptFu
func (s *rocksdb) GetRaw(ctx context.Context, col CF, key []byte, opts ...ReadOptFunc) (value []byte, err error) {
ro := &readOpts{}
ro.applyOptions(opts)
if ro.withNoMerge {
if ro.opt != nil || ro.withNoMerge {
return s.getRaw(ctx, col, key, ro.opt)
}
task := s.newReadTask(ctx)
task.typ = cfGetRaw
task.ro = ro
task.cf = []CF{col}
task.key = [][]byte{key}
@ -750,9 +756,6 @@ func (s *rocksdb) GetRaw(ctx context.Context, col CF, key []byte, opts ...ReadOp
func (s *rocksdb) MultiGet(ctx context.Context, col CF, keys [][]byte, opts ...ReadOptFunc) (values []ValueGetter, err error) {
ro := &readOpts{}
ro.applyOptions(opts)
if ro.withNoMerge {
return s.multiGet(ctx, col, keys, ro.opt)
}
task := s.newReadTask(ctx)
task.typ = multiGet
@ -775,13 +778,12 @@ func (s *rocksdb) MultiGet(ctx context.Context, col CF, keys [][]byte, opts ...R
func (s *rocksdb) SetRaw(ctx context.Context, col CF, key []byte, value []byte, opts ...WriteOptFunc) error {
wo := &writeOpts{}
wo.applyOptions(opts)
if wo.withNoMerge {
if wo.opt != nil || wo.withNoMerge {
return s.set(ctx, col, key, value, wo.opt)
}
task := s.newWriteTask(ctx)
task.typ = cfPutEvent
task.wo = wo
task.data = putData{
cf: col,
key: key,
@ -798,14 +800,13 @@ func (s *rocksdb) SetRaw(ctx context.Context, col CF, key []byte, value []byte,
func (s *rocksdb) Delete(ctx context.Context, col CF, key []byte, opts ...WriteOptFunc) error {
wo := &writeOpts{}
wo.applyOptions(opts)
if wo.withNoMerge {
if wo.opt != nil || wo.withNoMerge {
return s.delete(ctx, col, key, wo.opt)
}
task := s.newWriteTask(ctx)
task.typ = cfDeleteEvent
task.ctx = ctx
task.wo = wo
task.data = putData{
cf: col,
key: key,
@ -821,13 +822,12 @@ func (s *rocksdb) Delete(ctx context.Context, col CF, key []byte, opts ...WriteO
func (s *rocksdb) DeleteRange(ctx context.Context, col CF, start, end []byte, opts ...WriteOptFunc) error {
wo := &writeOpts{}
wo.applyOptions(opts)
if wo.withNoMerge {
if wo.opt != nil || wo.withNoMerge {
return s.deleteRange(ctx, col, start, end, wo.opt)
}
task := s.newWriteTask(ctx)
task.typ = cfRangeDeleteEvent
task.wo = wo
task.data = putData{
cf: col,
start: start,
@ -871,9 +871,6 @@ func (s *rocksdb) List(ctx context.Context, col CF, prefix []byte, marker []byte
func (s *rocksdb) Write(ctx context.Context, batch WriteBatch, opts ...WriteOptFunc) error {
wo := &writeOpts{}
wo.applyOptions(opts)
if wo.withNoMerge {
return s.write(ctx, batch, wo.opt)
}
task := s.newWriteTask(ctx)
task.typ = batchEvent
@ -890,9 +887,6 @@ func (s *rocksdb) Write(ctx context.Context, batch WriteBatch, opts ...WriteOptF
func (s *rocksdb) Read(ctx context.Context, cols []CF, keys [][]byte, opts ...ReadOptFunc) (values []ValueGetter, err error) {
ro := &readOpts{}
ro.applyOptions(opts)
if ro.withNoMerge {
return s.read(ctx, cols, keys, ro.opt)
}
task := s.newReadTask(ctx)
task.typ = read
@ -914,6 +908,9 @@ func (s *rocksdb) Read(ctx context.Context, cols []CF, keys [][]byte, opts ...Re
func (s *rocksdb) FlushCF(ctx context.Context, col CF) error {
cf := s.getColumnFamily(col)
s.lock.Lock()
defer s.lock.Unlock()
if err := s.db.FlushCF(s.fo, cf); err != nil {
s.handleError(ctx, err)
return err

View File

@ -58,7 +58,6 @@ func newEngine(ctx context.Context, opt *Option) (*testEg, error) {
_opt = new(Option)
}
_opt.CreateIfMissing = true
_opt.Sync = true
_opt.ReadConcurrency = 10
_opt.ReadQueueLen = 10
_opt.WriteConcurrency = 10
@ -161,6 +160,46 @@ func TestInstance_SetGetRaw(t *testing.T) {
require.Equal(t, ErrNotFound, err)
}
func TestPersistedRead(t *testing.T) {
ctx := context.TODO()
eg, err := newEngine(ctx, &Option{DisableWal: true})
require.NoError(t, err)
defer eg.close()
k1 := []byte("key1")
v1 := []byte("value1")
err = eg.engine.SetRaw(ctx, defaultCF, []byte("a"), []byte("b"))
require.Nil(t, err)
// write to wal
wo := eg.engine.NewWriteOption()
wo.DisableWAL(false)
defer wo.Close()
err = eg.engine.SetRaw(ctx, defaultCF, k1, v1, WithWriteOption(wo))
require.Nil(t, err)
v11, err := eg.engine.GetRaw(ctx, defaultCF, k1)
require.Nil(t, err)
require.Equal(t, v1, v11)
// persisted read
ro := eg.engine.NewReadOption()
ro.SetReadTier(ReadTierPersisted)
defer ro.Close()
_, err = eg.engine.GetRaw(ctx, defaultCF, k1, WithReadOption(ro))
require.Equal(t, ErrNotFound, err)
err = eg.engine.FlushCF(ctx, defaultCF)
require.Nil(t, err)
v11, err = eg.engine.GetRaw(ctx, defaultCF, k1, WithReadOption(ro))
require.Nil(t, err)
require.Equal(t, v1, v11)
}
func TestWriteRead(t *testing.T) {
ctx := context.TODO()
eg, err := newEngine(ctx, nil)
@ -177,6 +216,7 @@ func TestWriteRead(t *testing.T) {
values := make([][]byte, n)
batch := eg.engine.NewWriteBatch()
defer batch.Close()
for i := 0; i < n; i++ {
keyStr := []byte(fmt.Sprintf("k%d", i))
valStr := []byte(fmt.Sprintf("v%d", i))

View File

@ -275,11 +275,11 @@ func (s *service) loop(ctx context.Context) {
})
return true
})
for _, task := range tasks {
if err := s.executeShardTask(ctx, task); err != nil {
span.Errorf("execute shard task[%+v] failed: %s", task, errors.Detail(err))
continue
}
}
for _, task := range tasks {
if err := s.executeShardTask(ctx, task); err != nil {
span.Errorf("execute shard task[%+v] failed: %s", task, errors.Detail(err))
continue
}
}
case <-trashShardCheckTicker.C:
@ -364,7 +364,6 @@ func (s *service) executeShardTask(ctx context.Context, task clustermgr.ShardTas
_span.Errorf("shard do checkpoint task[%+v] failed: %s", task, errors.Detail(err))
}
})
default:
}
return nil

View File

@ -512,21 +512,17 @@ func (s *shard) Checkpoint(ctx context.Context) error {
}
defer s.shardState.prepRWCheckDone()
appliedIndex := (*shardSM)(s).getAppliedIndex()
// save applied index and shard's info
stableIndex := s.GetStableIndex()
flush := false
if appliedIndex != stableIndex {
flush = true
if err := s.SaveShardInfo(ctx, true, false); err != nil {
return errors.Info(err, "save shard info failed")
}
if err := s.SaveShardInfo(ctx, true, flush); err != nil {
if errors.Is(err, errShardStopWriting) {
span.Info("shard is stop writing by delete")
return nil
}
return errors.Info(err, "save shard into failed")
// get persisted info
info, err := s.getShardInfoFromPersistentTier(ctx)
if err != nil {
return errors.Info(err, "get shard info from persist layer failed")
}
appliedIndex := info.AppliedIndex
// truncate raft log finally
if appliedIndex > s.shardInfoMu.lastTruncatedIndex+s.cfg.TruncateWalLogInterval*2 {
@ -609,20 +605,15 @@ func (s *shard) SaveShardInfo(ctx context.Context, withLock bool, flush bool) er
s.shardInfoMu.Lock()
defer s.shardInfoMu.Unlock()
}
kvStore := s.store.KVStore()
key := s.shardKeys.encodeShardInfoKey()
value, err := s.shardInfoMu.shardInfo.Marshal()
if err != nil {
return err
}
if !flush {
wo := kvStore.NewWriteOption()
defer wo.Close()
return kvStore.SetRaw(ctx, dataCF, key, value, kvstore.WithWriteOption(wo))
return kvStore.SetRaw(ctx, dataCF, key, value)
}
if err := kvStore.SetRaw(ctx, dataCF, key, value); err != nil {
return err
}
@ -777,10 +768,6 @@ func (s *shard) GetAppliedIndex() uint64 {
return (*shardSM)(s).getAppliedIndex()
}
func (s *shard) GetStableIndex() uint64 {
return s.shardInfoMu.lastStableIndex
}
func (s *shard) GetSuid() proto.Suid {
return s.suid
}
@ -892,6 +879,24 @@ func (s *shard) getLeader(withLock bool) (clustermgr.ShardUnit, error) {
}, nil
}
func (s *shard) getShardInfoFromPersistentTier(ctx context.Context) (info clustermgr.Shard, err error) {
kvStore := s.store.KVStore()
key := s.shardKeys.encodeShardInfoKey()
ro := kvStore.NewReadOption()
ro.SetReadTier(kvstore.ReadTierPersisted)
defer ro.Close()
value, err := kvStore.GetRaw(ctx, dataCF, key, kvstore.WithReadOption(ro))
if err != nil {
return
}
if err = info.Unmarshal(value); err != nil {
return
}
return
}
func convertStoppingWriteErr(err error) error {
if errors.Is(err, errShardStopWriting) {
return apierr.ErrShardRouteVersionNeedUpdate

View File

@ -254,7 +254,7 @@ func TestServer_BlobList(t *testing.T) {
mockShard, shardClean := newMockShard(t)
defer shardClean()
err := mockShard.shard.SaveShardInfo(ctx, false, true)
err := mockShard.shard.SaveShardInfo(ctx, false, false)
require.Nil(t, err)
blobs := make([]cproto.Blob, 0)
@ -338,6 +338,18 @@ func TestServer_Snapshot(t *testing.T) {
require.Nil(t, err)
}
func TestShardInfo(t *testing.T) {
ctx := context.Background()
mockShard, shardClean := newMockShard(t)
defer shardClean()
err := mockShard.shard.SaveShardInfo(ctx, false, true)
require.Nil(t, err)
_, err = mockShard.shard.getShardInfoFromPersistentTier(ctx)
require.Nil(t, err)
}
func checkItemEqual(t *testing.T, shard *mockShard, id []byte, item *proto.Item) {
ret, err := shard.shard.GetItem(ctx, OpHeader{
ShardKeys: [][]byte{id},

View File

@ -123,6 +123,7 @@ func newMockShard(tb testing.TB) (*mockShard, func()) {
func TestServerShard_Checkpoint(t *testing.T) {
mockShard, shardClean := newMockShard(t)
defer shardClean()
mockShard.shard.SaveShardInfo(ctx, false, true)
gomock.InOrder(mockShard.mockRaftGroup.EXPECT().Truncate(A, A).AnyTimes().Return(nil))
err := mockShard.shard.Checkpoint(ctx)
require.Nil(t, err)
@ -213,7 +214,4 @@ func TestServerShard_Stats(t *testing.T) {
index := mockShard.shard.GetAppliedIndex()
require.Equal(t, uint64(0), index)
stableIdx := mockShard.shard.GetStableIndex()
require.Equal(t, uint64(0), stableIdx)
}