diff --git a/blobstore/blobnode/core/disk/disk.go b/blobstore/blobnode/core/disk/disk.go index 64a35b3dd..75021f94b 100644 --- a/blobstore/blobnode/core/disk/disk.go +++ b/blobstore/blobnode/core/disk/disk.go @@ -16,7 +16,6 @@ package disk import ( "context" - "fmt" "math/rand" "os" "path/filepath" @@ -462,15 +461,6 @@ func newDiskStorage(ctx context.Context, conf core.Config) (ds *DiskStorage, err conf.HandleIOError(context.Background(), dm.DiskID, err) }) - // io visualization: init meta io stat - metaios, err := flow.NewIOFlowStat(fmt.Sprintf("md_%v", dm.DiskID), conf.IOStatFileDryRun) - if err != nil { - span.Errorf("Failed new io flow stat, err:%v", err) - return nil, err - } - - sb.SetIOStat(metaios) - // io visualization: init data io stat dataios, err := flow.NewIOFlowStat(dm.DiskID.ToString(), conf.IOStatFileDryRun) if err != nil { diff --git a/blobstore/blobnode/core/disk/superblock.go b/blobstore/blobnode/core/disk/superblock.go index 4c44cd819..ceb9f91a9 100644 --- a/blobstore/blobnode/core/disk/superblock.go +++ b/blobstore/blobnode/core/disk/superblock.go @@ -23,7 +23,6 @@ import ( "strings" bnapi "github.com/cubefs/cubefs/blobstore/api/blobnode" - "github.com/cubefs/cubefs/blobstore/blobnode/base/flow" "github.com/cubefs/cubefs/blobstore/blobnode/core" "github.com/cubefs/cubefs/blobstore/blobnode/core/storage" "github.com/cubefs/cubefs/blobstore/blobnode/db" @@ -373,10 +372,6 @@ func (s *SuperBlock) SetHandlerIOError(handleIOError func(err error)) { s.db.SetHandleIOError(handleIOError) } -func (s *SuperBlock) SetIOStat(stat *flow.IOFlowStat) { - s.db.SetIOStat(stat) -} - func (s *SuperBlock) Close(ctx context.Context) error { s.db.Close(ctx) // db will be automatically closed when gc diff --git a/blobstore/blobnode/db/config.go b/blobstore/blobnode/db/config.go index 4d0509d58..0599b85eb 100644 --- a/blobstore/blobnode/db/config.go +++ b/blobstore/blobnode/db/config.go @@ -14,52 +14,20 @@ package db -import ( - "github.com/cubefs/cubefs/blobstore/blobnode/base/priority" - "github.com/cubefs/cubefs/blobstore/blobnode/base/qos" -) - const ( - DefaultBatchProcessCount = 1024 - DefaultWritePriRatio = 0.7 // 70 % - DefaultWriteBufferSize = 1 << 20 // 1 M - DefaultLRUCache = 256 << 20 // 256 M + DefaultLRUCache = 256 << 20 // 256 M ) type MetaConfig struct { - MetaRootPrefix string `json:"meta_root_prefix"` - SupportInline bool `json:"support_inline"` - TinyFileThresholdB int `json:"tinyfile_threshold_B"` - Sync bool `json:"sync"` - BatchProcessCount int64 `json:"batch_process_count"` - WritePriRatio float64 `json:"write_pri_ratio"` - MetaQos qos.LevelConfig `json:"meta_qos"` - LRUCacheSize uint64 `json:"cache_size"` + MetaRootPrefix string `json:"meta_root_prefix"` + SupportInline bool `json:"support_inline"` + TinyFileThresholdB int `json:"tinyfile_threshold_B"` + Sync bool `json:"sync"` + LRUCacheSize uint64 `json:"cache_size"` } func initConfig(conf *MetaConfig) error { // check params - for name := range conf.MetaQos { - if !priority.IsValidPriName(name) { - return ErrWrongConfig - } - } - - if conf.BatchProcessCount <= 0 { - conf.BatchProcessCount = DefaultBatchProcessCount - } - - // [ 55%, 95% ] - if conf.WritePriRatio <= 0 { - conf.WritePriRatio = DefaultWritePriRatio - } - if conf.WritePriRatio <= 0.55 { - conf.WritePriRatio = 0.55 - } - if conf.WritePriRatio >= 0.95 { - conf.WritePriRatio = 0.95 - } - if conf.LRUCacheSize == 0 { conf.LRUCacheSize = DefaultLRUCache } diff --git a/blobstore/blobnode/db/metadb.go b/blobstore/blobnode/db/metadb.go index 51fb081ad..75508270f 100644 --- a/blobstore/blobnode/db/metadb.go +++ b/blobstore/blobnode/db/metadb.go @@ -19,29 +19,13 @@ import ( "errors" "os" "sync" - "time" - "golang.org/x/time/rate" - - bnapi "github.com/cubefs/cubefs/blobstore/api/blobnode" bncom "github.com/cubefs/cubefs/blobstore/blobnode/base" - "github.com/cubefs/cubefs/blobstore/blobnode/base/flow" - "github.com/cubefs/cubefs/blobstore/blobnode/base/limitio" - pri "github.com/cubefs/cubefs/blobstore/blobnode/base/priority" bloberr "github.com/cubefs/cubefs/blobstore/common/errors" - "github.com/cubefs/cubefs/blobstore/common/iostat" rdb "github.com/cubefs/cubefs/blobstore/common/kvstore" "github.com/cubefs/cubefs/blobstore/common/trace" ) -const ( - _limited = "metalimited" -) - -const ( - _bufferDepth = 128 -) - var ( ErrStopped = errors.New("db: err stopped") ErrShardMetaNotDir = errors.New("db: shard meta not directory") @@ -55,31 +39,15 @@ type MetaHandler interface { DeleteRange(ctx context.Context, start, end []byte) error Flush(ctx context.Context) error NewIterator(ctx context.Context, opts ...rdb.OpOption) rdb.Iterator - SetIOStat(stat *flow.IOFlowStat) SetHandleIOError(handler func(err error)) Close(ctx context.Context) error } -type MetaDBWapper struct { - *metadb -} - type metadb struct { - lock sync.RWMutex - db rdb.KVStore - path string - - writeReqs chan Request - delReqs chan Request - config MetaConfig - - iostat *flow.IOFlowStat // io visualization - limiter []*rate.Limiter // io limiter - - closeCh chan struct{} - closed bool - - onClosed func() + once sync.Once + db rdb.KVStore + path string + config MetaConfig handleIOError func(err error) } @@ -93,30 +61,19 @@ func (md *metadb) SetHandleIOError(handler func(err error)) { md.handleIOError = handler } -func (md *metadb) SetIOStat(stat *flow.IOFlowStat) { - md.iostat = stat -} - func (md *metadb) Get(ctx context.Context, key []byte) (value []byte, err error) { - mgr, iot := md.getIoType(ctx) - - mgr.ReadBegin(1) - defer mgr.ReadEnd(time.Now()) - - md.applyToken(ctx, iot) - value, err = md.db.Get(key) md.handleError(err) return } -func (md *metadb) DirectPut(ctx context.Context, kv rdb.KV) (err error) { +func (md *metadb) Put(ctx context.Context, kv rdb.KV) (err error) { err = md.db.Put(kv) md.handleError(err) return } -func (md *metadb) DirectDelete(ctx context.Context, key []byte) (err error) { +func (md *metadb) Delete(ctx context.Context, key []byte) (err error) { err = md.db.Delete(key) md.handleError(err) return @@ -128,76 +85,10 @@ func (md *metadb) Flush(ctx context.Context) (err error) { return } -func (md *metadb) Put(ctx context.Context, kv rdb.KV) (err error) { - req := Request{ - Type: msgPut, - Data: kv, - } - - mgr, iot := md.getIoType(ctx) - - mgr.WriteBegin(1) - defer mgr.WriteEnd(time.Now()) - - resCh := req.Register() - - md.applyToken(ctx, iot) - md.writeReqs <- req - - select { - case <-md.closeCh: - return ErrStopped - case x := <-resCh: - return x.err - } -} - -func (md *metadb) Delete(ctx context.Context, key []byte) (err error) { - req := Request{ - Type: msgDel, - Data: rdb.KV{Key: key}, - } - - mgr, iot := md.getIoType(ctx) - - mgr.WriteBegin(1) - defer mgr.WriteEnd(time.Now()) - - resCh := req.Register() - - md.applyToken(ctx, iot) - md.delReqs <- req - - select { - case <-md.closeCh: - return ErrStopped - case x := <-resCh: - return x.err - } -} - func (md *metadb) DeleteRange(ctx context.Context, start, end []byte) (err error) { - req := Request{ - Type: msgDelRange, - Data: rdb.Range{Start: start, Limit: end}, - } - - mgr, iot := md.getIoType(ctx) - - mgr.WriteBegin(1) - defer mgr.WriteEnd(time.Now()) - - resCh := req.Register() - - md.applyToken(ctx, iot) - md.delReqs <- req - - select { - case <-md.closeCh: - return ErrStopped - case x := <-resCh: - return x.err - } + err = md.db.DeleteRange(start, end) + md.handleError(err) + return } func (md *metadb) NewIterator(ctx context.Context, opts ...rdb.OpOption) rdb.Iterator { @@ -207,163 +98,18 @@ func (md *metadb) NewIterator(ctx context.Context, opts ...rdb.OpOption) rdb.Ite func (md *metadb) Close(ctx context.Context) (err error) { span := trace.SpanFromContextSafe(ctx) - md.lock.Lock() - defer md.lock.Unlock() + md.once.Do(func() { + span.Infof("=== meta db:%v close ===", md.path) - span.Infof("=== meta db:%v close ===", md.path) - - if md.closed { - span.Panicf("can not happened. meta:%v", md.path) - return - } - - if md.onClosed != nil { - md.onClosed() - } - - if md.closeCh != nil { - close(md.closeCh) - } - - db := md.db - md.db = nil - - err = db.Close() - if err != nil { - span.Errorf("Failed close meta:%s, err:%v", md.path, err) - } - - md.closed = true + err = md.db.Close() + if err != nil { + span.Errorf("Failed close meta:%s, err:%v", md.path, err) + } + }) return err } -func (md *metadb) loopWorker() { - span, ctx := trace.StartSpanFromContextWithTraceID(context.Background(), "", "kvloop "+md.path) - - for { - var reqs []Request - - select { - case <-md.closeCh: - span.Warn("end the loop.") - return - case req := <-md.writeReqs: - reqs = append(reqs, req) - md.processRequests(ctx, reqs) - case req := <-md.delReqs: - reqs = append(reqs, req) - md.processRequests(ctx, reqs) - } - } -} - -func batchPopChannel(list []Request, reqChan chan Request, limit int) []Request { - if limit < 1 { - limit = 1 - } - for i := 0; i < limit; i++ { - var done bool - select { - case req := <-reqChan: - list = append(list, req) - default: - done = true - } - if done { - break - } - } - - return list -} - -func (md *metadb) processRequests(ctx context.Context, reqs []Request) { - span := trace.SpanFromContextSafe(ctx) - - batchCnt := int(md.config.BatchProcessCount) - writeCnt := int(float64(batchCnt) * md.config.WritePriRatio) - - limit := writeCnt - reqs = batchPopChannel(reqs, md.writeReqs, limit) - - limit = batchCnt - len(reqs) - reqs = batchPopChannel(reqs, md.delReqs, limit) - - if err := md.doBatch(ctx, reqs); err != nil { - span.Errorf("Failed doBatch reqs:%d, writeCnt:%d err:%v", len(reqs), writeCnt, err) - } -} - -func (md *metadb) doBatch(ctx context.Context, reqs []Request) (err error) { - writeBatch := md.db.NewWriteBatch() - defer writeBatch.Destroy() - - for _, req := range reqs { - switch req.Type { - case msgPut: - kv := req.Data.(rdb.KV) - writeBatch.Put(kv.Key, kv.Value) - case msgDel: - kv := req.Data.(rdb.KV) - writeBatch.Delete(kv.Key) - case msgDelRange: - r := req.Data.(rdb.Range) - writeBatch.DeleteRange(r.Start, r.Limit) - default: - } - } - - err = md.db.DoBatch(writeBatch) - if err != nil { - md.handleError(err) - } - - for _, req := range reqs { - req.Trigger(Result{err: err}) - } - - return err -} - -func (md *metadb) getIoType(ctx context.Context) (iostat.StatMgrAPI, bnapi.IOType) { - iot := bnapi.GetIoType(ctx) - mgr := md.iostat.GetStatMgr(iot) - return mgr, iot -} - -func (md *metadb) applyToken(ctx context.Context, iot bnapi.IOType) (n int64) { - priority := pri.GetPriority(iot) - - limiter := md.limiter[priority] - if limiter == nil { - return - } - now := time.Now() - reserve := limiter.ReserveN(now, 1) - delay := reserve.DelayFrom(now) - if delay == 0 { - return - } - t := time.NewTimer(delay) - defer t.Stop() - addTrackTag(ctx, _limited) - - select { - case <-t.C: - return - case <-ctx.Done(): - reserve.Cancel() - return - } -} - -func addTrackTag(ctx context.Context, name string) { - if open := limitio.IsLimitTrack(ctx); open { - limitio.AddTrackTag(ctx, name) - } -} - func newRocksDB(path string, conf MetaConfig) (db rdb.KVStore, err error) { if path == "" { return nil, bloberr.ErrInvalidParam @@ -378,7 +124,7 @@ func newRocksDB(path string, conf MetaConfig) (db rdb.KVStore, err error) { return nil, ErrShardMetaNotDir } - db, err = rdb.OpenDB(path, rdb.WithCatchSize(conf.LRUCacheSize)) + db, err = rdb.OpenDB(path, rdb.WithCatchSize(conf.LRUCacheSize), rdb.WithSyncMode(conf.Sync)) if err != nil { return } @@ -386,10 +132,10 @@ func newRocksDB(path string, conf MetaConfig) (db rdb.KVStore, err error) { return db, nil } -func newMetaDB(dirpath string, config MetaConfig) (md *metadb, err error) { - span, _ := trace.StartSpanFromContextWithTraceID(context.Background(), "", "NewKVDB "+dirpath) +func newMetaDB(path string, config MetaConfig) (md *metadb, err error) { + span, _ := trace.StartSpanFromContextWithTraceID(context.Background(), "", "NewKVDB "+path) - span.Infof("dirpath:%s, config:%v", dirpath, config) + span.Infof("path:%s, config:%v", path, config) err = initConfig(&config) if err != nil { @@ -397,48 +143,23 @@ func newMetaDB(dirpath string, config MetaConfig) (md *metadb, err error) { return } - rocksdb, err := newRocksDB(dirpath, config) + rocksdb, err := newRocksDB(path, config) if err != nil { span.Errorf("Failed New Rocksdb %v", err) return nil, err } md = &metadb{ - db: rocksdb, - path: dirpath, - writeReqs: make(chan Request, _bufferDepth), - delReqs: make(chan Request, _bufferDepth), - closeCh: make(chan struct{}), - config: config, + db: rocksdb, + path: path, + config: config, } - priLevels := pri.GetLevels() - md.limiter = make([]*rate.Limiter, len(priLevels)) - - for priority, name := range priLevels { - para, exist := config.MetaQos[name] - if !exist || para.Iops <= 0 { - // No flow control by default without configuration - continue - } - controller := rate.NewLimiter(rate.Limit(para.Iops), 2*int(para.Iops)) - md.limiter[priority] = controller - } - - span.Debugf("New KV(%s) DB(%v) success", dirpath, config) - - go md.loopWorker() + span.Debugf("New KV(%s) DB(%v) success", path, config) return md, nil } func NewMetaHandler(dirpath string, config MetaConfig) (mh MetaHandler, err error) { - md, err := newMetaDB(dirpath, config) - if err != nil { - return nil, err - } - - w := &MetaDBWapper{metadb: md} - - return w, nil + return newMetaDB(dirpath, config) } diff --git a/blobstore/blobnode/db/metadb_bench_test.go b/blobstore/blobnode/db/metadb_bench_test.go index 45fef081c..48942c059 100644 --- a/blobstore/blobnode/db/metadb_bench_test.go +++ b/blobstore/blobnode/db/metadb_bench_test.go @@ -104,7 +104,7 @@ func BenchmarkKVDB_DirectPut(b *testing.B) { Key: []byte(key), Value: value[:], } - err = md.DirectPut(ctx, kv) + err = md.Put(ctx, kv) require.NoError(b, err) }(i) } @@ -205,7 +205,7 @@ func BenchmarkKVDB_DirectDelete(b *testing.B) { Key: []byte(key), Value: value[:], } - err = md.DirectPut(ctx, kv) + err = md.Put(ctx, kv) require.NoError(b, err) }(i) } @@ -219,7 +219,7 @@ func BenchmarkKVDB_DirectDelete(b *testing.B) { go func(i int) { defer wg.Done() key := fmt.Sprintf("%s-%d", string(expectedKey), i) - err = md.DirectDelete(ctx, []byte(key)) + err = md.Delete(ctx, []byte(key)) require.NoError(b, err) }(i) } diff --git a/blobstore/blobnode/db/metadb_test.go b/blobstore/blobnode/db/metadb_test.go index a4c0a97ed..f0f235d1c 100644 --- a/blobstore/blobnode/db/metadb_test.go +++ b/blobstore/blobnode/db/metadb_test.go @@ -24,7 +24,6 @@ import ( "path/filepath" "sync" "testing" - "time" "github.com/stretchr/testify/require" @@ -70,14 +69,14 @@ func TestKVDB_DirectOP(t *testing.T) { ctx := context.Background() - err = md.DirectPut(ctx, kv) + err = md.Put(ctx, kv) require.NoError(t, err) value, err := md.Get(ctx, expectedKey) require.NoError(t, err) require.Equal(t, expectedValue, value) - err = md.DirectDelete(ctx, expectedKey) + err = md.Delete(ctx, expectedKey) require.NoError(t, err) _, err = md.Get(ctx, expectedKey) @@ -255,8 +254,6 @@ func TestKVDB_Close(t *testing.T) { require.NoError(t, err) defer os.RemoveAll(testDir) - span, _ := trace.StartSpanFromContextWithTraceID(context.Background(), "", "metadb") - diskdir := filepath.Join(testDir, "disk1/") err = os.MkdirAll(diskdir, 0o755) require.NoError(t, err) @@ -265,26 +262,9 @@ func TestKVDB_Close(t *testing.T) { require.NoError(t, err) require.NotNil(t, md) - var cnt int - done := make(chan struct{}) - - md.(*MetaDBWapper).onClosed = func() { - cnt++ - close(done) - } - // Trigger Close err = md.Close(context.Background()) require.NoError(t, err) - - select { - case <-done: - span.Infof("success gc") - case <-time.After(10 * time.Second): - t.Fail() - } - - require.Equal(t, 1, cnt) } func TestDeleteRange(t *testing.T) { diff --git a/blobstore/blobnode/db/request.go b/blobstore/blobnode/db/request.go deleted file mode 100644 index b87b31847..000000000 --- a/blobstore/blobnode/db/request.go +++ /dev/null @@ -1,43 +0,0 @@ -// Copyright 2022 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 db - -type MessageType int32 - -const ( - msgPut MessageType = iota + 1 - msgDel - msgDelRange -) - -type Request struct { - Type MessageType - Data interface{} - resCh chan Result -} - -type Result struct { - err error -} - -func (r *Request) Trigger(x Result) { - r.resCh <- x - close(r.resCh) -} - -func (r *Request) Register() <-chan Result { - r.resCh = make(chan Result, 1) - return r.resCh -} diff --git a/blobstore/blobnode/db/request_test.go b/blobstore/blobnode/db/request_test.go deleted file mode 100644 index 85c9ec0f4..000000000 --- a/blobstore/blobnode/db/request_test.go +++ /dev/null @@ -1,39 +0,0 @@ -// Copyright 2022 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 db - -import ( - "testing" - - "github.com/stretchr/testify/require" - - rdb "github.com/cubefs/cubefs/blobstore/common/kvstore" -) - -func TestRequest_Register(t *testing.T) { - req := Request{ - Type: msgPut, - Data: rdb.KV{Key: []byte{0x1}, Value: []byte{0x2}}, - } - - ch := req.Register() - - go func() { - req.Trigger(Result{err: ErrStopped}) - }() - - x := <-ch - require.Error(t, x.err) -} diff --git a/blobstore/common/kvstore/db.go b/blobstore/common/kvstore/db.go index 472403fb1..9b9169196 100644 --- a/blobstore/common/kvstore/db.go +++ b/blobstore/common/kvstore/db.go @@ -467,6 +467,16 @@ func (s *instance) Delete(key []byte) (err error) { return <-task.err } +func (s *instance) DeleteRange(start, end []byte) error { + task := &writeTask{ + typ: rangeDeleteEvent, + data: rangeDelete{start, end}, + err: make(chan error, 1), + } + s.wchan <- task + return <-task.err +} + func (s *instance) DeleteBatch(keys [][]byte, safe bool) (err error) { b := []batch{} for _, key := range keys { diff --git a/blobstore/common/kvstore/proto.go b/blobstore/common/kvstore/proto.go index 1266f0a1a..e007124ca 100644 --- a/blobstore/common/kvstore/proto.go +++ b/blobstore/common/kvstore/proto.go @@ -18,6 +18,7 @@ type KVStorage interface { Put(kv KV) error Get(key []byte) ([]byte, error) Delete(key []byte) error + DeleteRange(start, end []byte) error NewWriteBatch() *WriteBatch DeleteBatch(keys [][]byte, safe bool) error WriteBatch(kvs []KV, safe bool) error diff --git a/blobstore/common/kvstore/table.go b/blobstore/common/kvstore/table.go index d1339a90d..0cfd56618 100644 --- a/blobstore/common/kvstore/table.go +++ b/blobstore/common/kvstore/table.go @@ -82,6 +82,16 @@ func (t *table) Delete(key []byte) (err error) { return <-task.err } +func (t *table) DeleteRange(start, end []byte) (err error) { + task := &writeTask{ + typ: cfRangeDeleteEvent, + data: cfRangeDelete{t.cf, start, end}, + err: make(chan error, 1), + } + t.ins.wchan <- task + return <-task.err +} + func (t *table) Name() string { return t.name }