feat(blobnode): start mode split blobnode service

with: #1000159741

Signed-off-by: mawei029 <mawei2@oppo.com>
This commit is contained in:
mawei029 2025-06-10 16:56:29 +08:00 committed by slasher
parent 3197d47f54
commit 83f0b7c140
6 changed files with 147 additions and 31 deletions

View File

@ -17,7 +17,6 @@ package blobnode
import (
"context"
"net/http"
"os"
"strconv"
bnapi "github.com/cubefs/cubefs/blobstore/api/blobnode"
@ -47,6 +46,8 @@ const (
defaultInspectIntervalSec = 24 * 60 * 60 // 24 hour
defaultInspectRate = 4 * 1024 * 1024 // rate limit 4MB per second
defaultInspectLogChunkSize = uint(29)
defaultServiceBothBlobNodeWorker = "both_blobnode_worker"
)
var (
@ -55,6 +56,8 @@ var (
ErrValueOutOfLimit = errors.New("value out of limit")
)
type StartMode string
type Config struct {
cmd.Config
core.HostInfo
@ -76,12 +79,16 @@ type Config struct {
DeleteQpsLimitPerDisk int `json:"delete_qps_limit_per_disk"`
InspectConf DataInspectConf `json:"inspect_conf"`
StartMode StartMode `json:"start_mode"`
}
func configInit(config *Config) {
if config.StartMode == proto.ServiceNameWorker {
return
}
if len(config.Disks) == 0 {
log.Fatalf("disk list is empty")
os.Exit(1)
}
if config.HeartbeatIntervalSec <= 0 {
@ -116,6 +123,12 @@ func configInit(config *Config) {
defaulter.LessOrEqual(&config.InspectConf.Record.ChunkBits, defaultInspectLogChunkSize)
defaulter.LessOrEqual(&config.HostInfo.DiskType, proto.DiskTypeHDD)
defaulter.Empty((*string)(&config.StartMode), defaultServiceBothBlobNodeWorker)
// Start the worker process service by the ${StartMode}: only worker, only blobnode, both
if config.StartMode != defaultServiceBothBlobNodeWorker && config.StartMode != proto.ServiceNameWorker && config.StartMode != proto.ServiceNameBlobNode {
log.Fatalf("fail to start service, mode is %s, please use valid start mode", config.StartMode)
}
}
func (s *Service) changeLimit(ctx context.Context, c Config) {

View File

@ -301,11 +301,58 @@ func (s *Service) fixDiskConf(config *core.Config) {
}
func NewService(conf Config) (svr *Service, err error) {
span, ctx := trace.StartSpanFromContext(context.Background(), "NewBlobNodeService")
_, ctx := trace.StartSpanFromContext(context.Background(), "NewBlobNodeService")
configInit(&conf)
clusterMgrCli := cmapi.New(conf.Clustermgr)
svr = &Service{
ClusterMgrClient: clusterMgrCli,
Disks: make(map[proto.DiskID]core.DiskAPI),
Conf: &conf,
closeCh: make(chan struct{}),
}
// start worker service
if conf.StartMode == proto.ServiceNameWorker || conf.StartMode == defaultServiceBothBlobNodeWorker {
startWorkerService(ctx, svr, conf, clusterMgrCli)
}
// start blobndoe service
if conf.StartMode == proto.ServiceNameBlobNode || conf.StartMode == defaultServiceBothBlobNodeWorker {
err = startBlobnodeService(ctx, svr, conf, clusterMgrCli)
}
return
}
func startWorkerService(ctx context.Context, svr *Service, conf Config, clusterMgrCli *cmapi.Client) {
span := trace.SpanFromContextSafe(ctx)
span.Debug("start worker service...")
node := cmapi.ServiceNode{
ClusterID: uint64(conf.ClusterID),
Name: proto.ServiceNameWorker,
Host: conf.Host,
Idc: conf.IDC,
}
err := clusterMgrCli.RegisterService(ctx, node, TickInterval, HeartbeatTicks, ExpiresTicks)
if err != nil {
span.Fatalf("worker register to clusterMgr error:%+v", err)
}
svr.WorkerService, err = NewWorkerService(&conf.WorkerConfig, clusterMgrCli, conf.ClusterID, conf.IDC)
if err != nil {
span.Fatalf("Failed to new worker service, err: %v", err)
}
}
func startBlobnodeService(ctx context.Context, svr *Service, conf Config, clusterMgrCli *cmapi.Client) (err error) {
span := trace.SpanFromContextSafe(ctx)
span.Debug("start blobnode service...")
node := cmapi.ServiceNode{
ClusterID: uint64(conf.ClusterID),
Name: proto.ServiceNameBlobNode,
@ -328,42 +375,34 @@ func NewService(conf Config) (svr *Service, err error) {
if err = registerNode(ctx, clusterMgrCli, &conf); err != nil {
span.Fatalf("fail to register node to clusterMgr, err:%+v", err)
return nil, err
}
registeredDisks, err := clusterMgrCli.ListHostDisk(ctx, conf.Host)
if err != nil {
span.Errorf("Failed ListDisk from clusterMgr. err:%+v", err)
return nil, err
return err
}
span.Infof("registered disks: %v", registeredDisks)
check := isAllInConfig(ctx, registeredDisks, &conf)
if !check {
span.Errorf("no all registered normal disk in config")
return nil, errors.New("registered disk not in config")
return errors.New("registered disk not in config")
}
span.Infof("registered disks are all in config")
svr = &Service{
ClusterMgrClient: clusterMgrCli,
Disks: make(map[proto.DiskID]core.DiskAPI),
Conf: &conf,
DeleteQpsLimitPerDisk: keycount.New(conf.DeleteQpsLimitPerDisk),
DeleteQpsLimitPerKey: keycount.NewBlockingKeyCountLimit(1),
ChunkLimitPerVuid: keycount.New(1),
DiskLimitRegister: keycount.New(1),
InspectLimiterPerKey: keycount.New(1),
BrokenLimitPerDisk: keycount.New(1),
closeCh: make(chan struct{}),
}
svr.Conf = &conf
svr.DeleteQpsLimitPerDisk = keycount.New(conf.DeleteQpsLimitPerDisk)
svr.DeleteQpsLimitPerKey = keycount.NewBlockingKeyCountLimit(1)
svr.ChunkLimitPerVuid = keycount.New(1)
svr.DiskLimitRegister = keycount.New(1)
svr.InspectLimiterPerKey = keycount.New(1)
svr.BrokenLimitPerDisk = keycount.New(1)
switchMgr := taskswitch.NewSwitchMgr(clusterMgrCli)
svr.inspectMgr, err = NewDataInspectMgr(svr, conf.InspectConf, switchMgr)
if err != nil {
return nil, err
return err
}
svr.ctx, svr.cancel = context.WithCancel(context.Background())
@ -439,7 +478,7 @@ func NewService(conf Config) (svr *Service, err error) {
if err = setDefaultIOStat(conf.DiskConfig.IOStatFileDryRun); err != nil {
span.Errorf("Failed set default iostat file, err:%v", err)
return nil, err
return err
}
callBackFn := func(conf []byte) error {
@ -457,12 +496,6 @@ func NewService(conf Config) (svr *Service, err error) {
}
config.Register(callBackFn)
svr.WorkerService, err = NewWorkerService(&conf.WorkerConfig, clusterMgrCli, conf.ClusterID, conf.IDC)
if err != nil {
span.Errorf("Failed to new worker service, err: %v", err)
return
}
// background loop goroutines
go svr.loopHeartbeatToClusterMgr()
go svr.loopReportChunkInfoToClusterMgr()

View File

@ -244,8 +244,8 @@ func TestHandleDiskDrop(t *testing.T) {
require.Equal(t, di.Status, proto.DiskStatusNormal)
}
func TestService2(t *testing.T) {
workDir, err := os.MkdirTemp(os.TempDir(), defaultSvrTestDir+"Service2")
func TestServiceError(t *testing.T) {
workDir, err := os.MkdirTemp(os.TempDir(), defaultSvrTestDir+"ServiceError")
require.NoError(t, err)
defer os.Remove(workDir)
@ -1164,3 +1164,71 @@ func TestService_RegisterNode(t *testing.T) {
err = registerNode(ctx, svr.ClusterMgrClient, svr.Conf)
require.NotNil(t, err)
}
func TestService_OnlyWorker(t *testing.T) {
mcm := mockClusterMgr{
reqIdx: _mockDiskIdBase,
disks: []mockDiskInfo{},
}
mcmURL := runMockClusterMgr(&mcm)
cc := &cmapi.Config{}
cc.Hosts = []string{mcmURL}
conf := Config{
Clustermgr: cc,
StartMode: proto.ServiceNameWorker,
}
_, err := NewService(conf)
require.NoError(t, err)
}
func TestService_OnlyBlobnode(t *testing.T) {
workDir, err := os.MkdirTemp(os.TempDir(), defaultSvrTestDir+"OnlyBlobnode")
require.NoError(t, err)
defer os.RemoveAll(workDir)
path1 := filepath.Join(workDir, "disk1")
path2 := filepath.Join(workDir, "disk2")
mcm := mockClusterMgr{
reqIdx: _mockDiskIdBase,
disks: []mockDiskInfo{},
}
mcmURL := runMockClusterMgr(&mcm)
cc := &cmapi.Config{}
cc.Hosts = []string{mcmURL}
err = os.MkdirAll(workDir, 0o755)
require.NoError(t, err)
// must create meta dir
err = os.MkdirAll(core.GetMetaPath(path1, ""), 0o755)
require.NoError(t, err)
err = os.MkdirAll(core.GetMetaPath(path2, ""), 0o755)
require.NoError(t, err)
conf := Config{
HostInfo: core.HostInfo{
IDC: "testIdc",
Rack: "testRack",
DiskType: proto.DiskTypeHDD,
},
Disks: []core.Config{
{BaseConfig: core.BaseConfig{Path: path1, AutoFormat: true, MaxChunks: 700}, MetaConfig: db.MetaConfig{}},
{BaseConfig: core.BaseConfig{Path: path2, AutoFormat: true, MaxChunks: 700}, MetaConfig: db.MetaConfig{}},
},
DiskConfig: core.RuntimeConfig{DiskReservedSpaceB: 1, CompactReservedSpaceB: 1},
Clustermgr: cc,
HeartbeatIntervalSec: 600,
InspectConf: DataInspectConf{Record: recordlog.Config{Dir: filepath.Join(workDir, "inspect")}},
StartMode: proto.ServiceNameBlobNode,
}
_, err = NewService(conf)
require.NoError(t, err)
}

View File

@ -31,6 +31,7 @@ func cmdGetService(c *grumble.Context) error {
proto.ServiceNameProxy,
proto.ServiceNameBlobNode,
proto.ServiceNameScheduler,
proto.ServiceNameWorker,
}
name := c.Args.String("name")
if name != "" {

View File

@ -28,6 +28,7 @@ const (
ServiceNameProxy = "PROXY"
ServiceNameScheduler = "SCHEDULER"
ServiceNameShardNode = "SHARDNODE"
ServiceNameWorker = "WORKER"
)
type (

View File

@ -134,7 +134,7 @@ func NewShardRepairMgr(
}
workerSelector := selector.MakeSelector(60*1000, func() (hosts []string, err error) {
return clusterMgrCli.GetService(context.Background(), proto.ServiceNameBlobNode, cfg.ClusterID)
return clusterMgrCli.GetService(context.Background(), proto.ServiceNameWorker, cfg.ClusterID)
})
failMsgSender, err := base.NewMsgSender(cfg.failedProducerConfig())
if err != nil {