refactor(blobnode): optimize the implementation of blobnode meta

Signed-off-by: yuxiaobo <yuxiaobo@oppo.com>
This commit is contained in:
yuxiaobo 2023-09-06 10:50:46 +08:00 committed by mujita
parent 2adc410ec6
commit 002bc4fbfb
11 changed files with 58 additions and 465 deletions

View File

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

View File

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

View File

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

View File

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

View File

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

View File

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

View File

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

View File

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

View File

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

View File

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

View File

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