feat(shardnode): inspect and clean up trash shards

with #22963391

Signed-off-by: xiejian <xiejian3@oppo.com>
This commit is contained in:
xiejian 2025-01-22 14:55:52 +08:00 committed by slasher
parent 828a7be5ac
commit 1cdfd8f95a
8 changed files with 158 additions and 6 deletions

View File

@ -119,6 +119,7 @@ const (
ShardTaskTypeClearShard = ShardTaskType(iota + 1)
ShardTaskTypeCheckpoint
ShardTaskTypeSyncRouteVersion
ShardTaskTypeCheckAndClear
)
type ShardUnitStatus uint8

View File

@ -333,6 +333,21 @@ func (mr *MockTransportMockRecorder) ShardReport(ctx, reports interface{}) *gomo
return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "ShardReport", reflect.TypeOf((*MockTransport)(nil).ShardReport), ctx, reports)
}
// ShardStats mocks base method.
func (m *MockTransport) ShardStats(ctx context.Context, host string, args shardnode.GetShardArgs) (shardnode.ShardStats, error) {
m.ctrl.T.Helper()
ret := m.ctrl.Call(m, "ShardStats", ctx, host, args)
ret0, _ := ret[0].(shardnode.ShardStats)
ret1, _ := ret[1].(error)
return ret0, ret1
}
// ShardStats indicates an expected call of ShardStats.
func (mr *MockTransportMockRecorder) ShardStats(ctx, host, args interface{}) *gomock.Call {
mr.mock.ctrl.T.Helper()
return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "ShardStats", reflect.TypeOf((*MockTransport)(nil).ShardStats), ctx, host, args)
}
// UpdateShard mocks base method.
func (m *MockTransport) UpdateShard(ctx context.Context, host string, args shardnode.UpdateShardArgs) error {
m.ctrl.T.Helper()
@ -688,6 +703,21 @@ func (mr *MockShardTransportMockRecorder) ResolveRaftAddr(ctx, diskID interface{
return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "ResolveRaftAddr", reflect.TypeOf((*MockShardTransport)(nil).ResolveRaftAddr), ctx, diskID)
}
// ShardStats mocks base method.
func (m *MockShardTransport) ShardStats(ctx context.Context, host string, args shardnode.GetShardArgs) (shardnode.ShardStats, error) {
m.ctrl.T.Helper()
ret := m.ctrl.Call(m, "ShardStats", ctx, host, args)
ret0, _ := ret[0].(shardnode.ShardStats)
ret1, _ := ret[1].(error)
return ret0, ret1
}
// ShardStats indicates an expected call of ShardStats.
func (mr *MockShardTransportMockRecorder) ShardStats(ctx, host, args interface{}) *gomock.Call {
mr.mock.ctrl.T.Helper()
return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "ShardStats", reflect.TypeOf((*MockShardTransport)(nil).ShardStats), ctx, host, args)
}
// UpdateShard mocks base method.
func (m *MockShardTransport) UpdateShard(ctx context.Context, host string, args shardnode.UpdateShardArgs) error {
m.ctrl.T.Helper()

View File

@ -67,6 +67,7 @@ type (
ResolveRaftAddr(ctx context.Context, diskID proto.DiskID) (string, error)
ResolveNodeAddr(ctx context.Context, diskID proto.DiskID) (string, error)
UpdateShard(ctx context.Context, host string, args shardnodeapi.UpdateShardArgs) error
ShardStats(ctx context.Context, host string, args shardnodeapi.GetShardArgs) (shardnodeapi.ShardStats, error)
}
)
@ -295,3 +296,7 @@ func (t *transport) ResolveNodeAddr(ctx context.Context, diskID proto.DiskID) (s
func (t *transport) UpdateShard(ctx context.Context, host string, args shardnodeapi.UpdateShardArgs) error {
return t.snClient.UpdateShard(ctx, host, args)
}
func (t *transport) ShardStats(ctx context.Context, host string, args shardnodeapi.GetShardArgs) (shardnodeapi.ShardStats, error) {
return t.snClient.GetShardStats(ctx, host, args)
}

View File

@ -279,6 +279,20 @@ func (m *MockSpaceShardHandler) EXPECT() *MockSpaceShardHandlerMockRecorder {
return m.recorder
}
// CheckAndClearShard mocks base method.
func (m *MockSpaceShardHandler) CheckAndClearShard(ctx context.Context) error {
m.ctrl.T.Helper()
ret := m.ctrl.Call(m, "CheckAndClearShard", ctx)
ret0, _ := ret[0].(error)
return ret0
}
// CheckAndClearShard indicates an expected call of CheckAndClearShard.
func (mr *MockSpaceShardHandlerMockRecorder) CheckAndClearShard(ctx interface{}) *gomock.Call {
mr.mock.ctrl.T.Helper()
return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "CheckAndClearShard", reflect.TypeOf((*MockSpaceShardHandler)(nil).CheckAndClearShard), ctx)
}
// Checkpoint mocks base method.
func (m *MockSpaceShardHandler) Checkpoint(ctx context.Context) error {
m.ctrl.T.Helper()

View File

@ -180,6 +180,7 @@ func (s *service) loop(ctx context.Context) {
reportTicker := time.NewTicker(time.Duration(s.cfg.ReportIntervalS) * time.Second)
routeUpdateTicker := time.NewTicker(time.Duration(s.cfg.RouteUpdateIntervalS) * time.Second)
checkpointTicker := time.NewTicker(time.Duration(s.cfg.CheckPointIntervalM) * time.Minute)
trashShardCheckTicker := time.NewTicker(time.Duration(s.cfg.ShardCheckAndClearIntervalH) * time.Hour)
defer func() {
heartbeatTicker.Stop()
@ -192,6 +193,7 @@ func (s *service) loop(ctx context.Context) {
diskReports := make([]clustermgr.ShardNodeDiskHeartbeatInfo, 0)
shardReports := make([]clustermgr.ShardUnitInfo, 0, 1<<10)
shards := make([]storage.ShardHandler, 0, 1<<10)
tasks := make([]clustermgr.ShardTask, 0, 1<<10)
for {
select {
@ -263,20 +265,44 @@ func (s *service) loop(ctx context.Context) {
case <-checkpointTicker.C:
span, ctx = trace.StartSpanFromContext(ctx, "do checkpoint")
disks := s.getAllDisks()
tasks = tasks[:0]
for _, disk := range disks {
shards = shards[:0]
disk.RangeShard(func(s storage.ShardHandler) bool {
shards = append(shards, s)
tasks = append(tasks, clustermgr.ShardTask{
TaskType: proto.ShardTaskTypeCheckpoint,
Suid: s.GetSuid(),
DiskID: disk.DiskID(),
})
return true
})
for _, shard := range shards {
err := shard.Checkpoint(ctx)
if err != nil {
span.Errorf("do checkpoint failed: %s", errors.Detail(err))
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:
span, ctx = trace.StartSpanFromContext(ctx, "trash shard check")
disks := s.getAllDisks()
tasks = tasks[:0]
for _, disk := range disks {
disk.RangeShard(func(s storage.ShardHandler) bool {
tasks = append(tasks, clustermgr.ShardTask{
TaskType: proto.ShardTaskTypeCheckAndClear,
Suid: s.GetSuid(),
DiskID: disk.DiskID(),
})
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
}
}
case <-s.closer.Done():
return
}
@ -325,6 +351,21 @@ func (s *service) executeShardTask(ctx context.Context, task clustermgr.ShardTas
curVersion, task.OldRouteVersion, task.RouteVersion)
}
})
case proto.ShardTaskTypeCheckAndClear:
s.taskPool.Run(func() {
_span, _ctx := trace.StartSpanFromContextWithTraceID(ctx, "", "shard-check-"+task.Suid.ToString())
if err := shard.CheckAndClearShard(_ctx); err != nil {
_span.Errorf("check trash shard task[%+v] failed: %s", task, errors.Detail(err))
}
})
case proto.ShardTaskTypeCheckpoint:
s.taskPool.Run(func() {
_span, _ctx := trace.StartSpanFromContextWithTraceID(ctx, "", "checkpoint-"+task.Suid.ToString())
if err := shard.Checkpoint(_ctx); err != nil {
_span.Errorf("shard do checkpoint task[%+v] failed: %s", task, errors.Detail(err))
}
})
default:
}
return nil

View File

@ -373,4 +373,5 @@ func initServiceConfig(cfg *Config) {
defaulter.LessOrEqual(&cfg.CheckPointIntervalM, int64(1))
defaulter.LessOrEqual(&cfg.WaitRepairCloseDiskIntervalS, int64(30))
defaulter.LessOrEqual(&cfg.WaitReOpenDiskIntervalS, int64(30))
defaulter.LessOrEqual(&cfg.ShardCheckAndClearIntervalH, int64(24))
}

View File

@ -27,6 +27,7 @@ import (
kvstore "github.com/cubefs/cubefs/blobstore/common/kvstorev2"
"github.com/cubefs/cubefs/blobstore/common/proto"
"github.com/cubefs/cubefs/blobstore/common/raft"
"github.com/cubefs/cubefs/blobstore/common/rpc"
"github.com/cubefs/cubefs/blobstore/common/sharding"
"github.com/cubefs/cubefs/blobstore/common/trace"
"github.com/cubefs/cubefs/blobstore/shardnode/base"
@ -73,6 +74,7 @@ type (
Stats(ctx context.Context) (shardnode.ShardStats, error)
GetSuid() proto.Suid
GetUnits() []clustermgr.ShardUnit
CheckAndClearShard(ctx context.Context) error
}
OpHeader struct {
RouteVersion proto.RouteVersion
@ -462,7 +464,8 @@ func (s *shard) Stats(ctx context.Context) (shardnode.ShardStats, error) {
defer s.shardState.prepRWCheckDone()
s.shardInfoMu.RLock()
units := s.shardInfoMu.Units
units := make([]clustermgr.ShardUnit, len(s.shardInfoMu.Units))
copy(units, s.shardInfoMu.Units)
routeVersion := s.shardInfoMu.RouteVersion
appliedIndex := s.shardInfoMu.AppliedIndex
rg := s.shardInfoMu.Range
@ -696,6 +699,62 @@ func (s *shard) DeleteShard(ctx context.Context, nodeHost string, clearData bool
return kvStore.FlushCF(ctx, dataCF)
}
func (s *shard) CheckAndClearShard(ctx context.Context) error {
span := trace.SpanFromContext(ctx)
stats, err := s.Stats(ctx)
if err != nil {
return err
}
if stats.LeaderDiskID != s.diskID {
return nil
}
raftStats := make(map[uint64]bool)
for _, p := range stats.RaftStat.Peers {
raftStats[p.NodeID] = p.RecentActive
}
for _, u := range stats.Units {
if u.Suid == s.GetSuid() {
continue
}
var active, ok bool
if active, ok = raftStats[uint64(u.DiskID)]; !ok {
return errors.Newf("disk[%d] not found in raft peers", u.DiskID)
}
if active {
continue
}
host, err := s.cfg.Transport.ResolveNodeAddr(ctx, u.DiskID)
if err != nil {
return errors.Info(err, "resolve node address failed", u.DiskID)
}
_, err = s.cfg.Transport.ShardStats(ctx, host, shardnodeapi.GetShardArgs{
DiskID: u.DiskID,
Suid: u.Suid,
})
if err == nil {
continue
}
if rpc.DetectStatusCode(err) != apierr.CodeShardNodeDiskNotFound {
return errors.Info(err, "get shard stats failed", u.DiskID, u.Suid)
}
err = s.UpdateShard(ctx, proto.ShardUpdateTypeRemoveMember, clustermgr.ShardUnit{
Suid: u.Suid,
DiskID: u.DiskID,
}, host)
if err != nil {
return errors.Info(err, "update shard failed", s.suid, u.DiskID)
}
span.Infof("remove shard[%d] suid[%d] from unit[%+v] done", s.suid.ShardID(), s.suid, u)
return err
}
return nil
}
func (s *shard) Start() {
// Do nothing because need to be improved later
}

View File

@ -71,6 +71,7 @@ type Config struct {
CheckPointIntervalM int64 `json:"check_point_interval_m"`
WaitRepairCloseDiskIntervalS int64 `json:"wait_repair_close_disk_interval_s"`
WaitReOpenDiskIntervalS int64 `json:"wait_re_open_disk_interval_s"`
ShardCheckAndClearIntervalH int64 `json:"shard_check_and_clear_interval_h"`
}
func newService(cfg *Config) *service {