diff --git a/blobstore/common/kvstorev2/kvstore.go b/blobstore/common/kvstorev2/kvstore.go index e5dd61b76..44b9d556e 100644 --- a/blobstore/common/kvstorev2/kvstore.go +++ b/blobstore/common/kvstorev2/kvstore.go @@ -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) diff --git a/blobstore/common/kvstorev2/rocksdb.go b/blobstore/common/kvstorev2/rocksdb.go index a670c89d2..db48f222d 100644 --- a/blobstore/common/kvstorev2/rocksdb.go +++ b/blobstore/common/kvstorev2/rocksdb.go @@ -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 diff --git a/blobstore/common/kvstorev2/rocksdb_test.go b/blobstore/common/kvstorev2/rocksdb_test.go index b7560d39c..4c4909ee5 100644 --- a/blobstore/common/kvstorev2/rocksdb_test.go +++ b/blobstore/common/kvstorev2/rocksdb_test.go @@ -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)) diff --git a/blobstore/shardnode/shard.go b/blobstore/shardnode/shard.go index 4b92bb54b..d89e0fb8a 100644 --- a/blobstore/shardnode/shard.go +++ b/blobstore/shardnode/shard.go @@ -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 diff --git a/blobstore/shardnode/storage/shard.go b/blobstore/shardnode/storage/shard.go index 002ecf691..d7a23bd82 100644 --- a/blobstore/shardnode/storage/shard.go +++ b/blobstore/shardnode/storage/shard.go @@ -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 diff --git a/blobstore/shardnode/storage/shard_sm_test.go b/blobstore/shardnode/storage/shard_sm_test.go index d29f972cc..7ecd291e7 100644 --- a/blobstore/shardnode/storage/shard_sm_test.go +++ b/blobstore/shardnode/storage/shard_sm_test.go @@ -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}, diff --git a/blobstore/shardnode/storage/shard_test.go b/blobstore/shardnode/storage/shard_test.go index 350a522a1..a439add59 100644 --- a/blobstore/shardnode/storage/shard_test.go +++ b/blobstore/shardnode/storage/shard_test.go @@ -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) }