feat(shardnode): transfer leader when handleEIO

with #22357426

Signed-off-by: xiejian <xiejian3@oppo.com>
This commit is contained in:
xiejian 2024-12-31 10:14:33 +08:00 committed by slasher
parent 759217aeab
commit 28c197a0e2
4 changed files with 84 additions and 0 deletions

View File

@ -453,6 +453,20 @@ func (mr *MockSpaceShardHandlerMockRecorder) TransferLeader(ctx, diskID interfac
return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "TransferLeader", reflect.TypeOf((*MockSpaceShardHandler)(nil).TransferLeader), ctx, diskID)
}
// TryTransferLeader mocks base method.
func (m *MockSpaceShardHandler) TryTransferLeader(ctx context.Context) error {
m.ctrl.T.Helper()
ret := m.ctrl.Call(m, "TryTransferLeader", ctx)
ret0, _ := ret[0].(error)
return ret0
}
// TryTransferLeader indicates an expected call of TryTransferLeader.
func (mr *MockSpaceShardHandlerMockRecorder) TryTransferLeader(ctx interface{}) *gomock.Call {
mr.mock.ctrl.T.Helper()
return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "TryTransferLeader", reflect.TypeOf((*MockSpaceShardHandler)(nil).TryTransferLeader), ctx)
}
// Update mocks base method.
func (m *MockSpaceShardHandler) Update(ctx context.Context, h storage.OpHeader, kv *storage.KV) error {
m.ctrl.T.Helper()

View File

@ -16,6 +16,7 @@
package shardnode
import (
"container/list"
"context"
"fmt"
"sync"
@ -185,6 +186,24 @@ func (s *service) handleEIO(ctx context.Context, diskID proto.DiskID, err error)
time.Sleep(5 * time.Second)
}
go func() {
failedShards := list.New()
disk.RangeShard(func(s storage.ShardHandler) bool {
failedShards.PushBack(s)
return true
})
for failedShards.Len() > 0 {
e := failedShards.Front()
failedShards.Remove(e)
s := e.Value.(storage.ShardHandler)
if err := s.TryTransferLeader(ctx); err != nil {
failedShards.PushBack(e)
}
}
span.Infof("disk[%d] transfer shards leader done", diskID)
}()
// wait for disk repairing
go s.waitRepairCloseDisk(ctx, disk)

View File

@ -204,6 +204,19 @@ func TestServerDisk_Raft(t *testing.T) {
}
t.Logf("add shard[%d] success", suid4)
// transfer leader
err = shardLeader.TryTransferLeader(ctx)
require.Nil(t, err)
for {
shard1, _err := disk1.GetShard(suid1)
require.Nil(t, _err)
stat, _ := shard1.Stats(ctx)
if stat.LeaderDiskID != proto.InvalidDiskID && stat.LeaderDiskID != leaderDiskID {
leaderDiskID = stat.LeaderDiskID
break
}
}
version += 1
delIdx := 0
for i, d := range disks {

View File

@ -69,6 +69,7 @@ type (
ShardItemHandler
GetRouteVersion() proto.RouteVersion
TransferLeader(ctx context.Context, diskID proto.DiskID) error
TryTransferLeader(ctx context.Context) (err error)
Checkpoint(ctx context.Context) error
Stats(ctx context.Context) (shardnode.ShardStats, error)
GetSuid() proto.Suid
@ -583,6 +584,43 @@ func (s *shard) TransferLeader(ctx context.Context, diskID proto.DiskID) error {
return s.raftGroup.LeaderTransfer(ctx, uint64(diskID))
}
func (s *shard) TryTransferLeader(ctx context.Context) (err error) {
span := trace.SpanFromContextSafe(ctx)
defer func() {
if errors.Is(err, apierr.ErrShardRouteVersionNeedUpdate) {
span.Debugf("shard[%d] suid[%d] is closed", s.suid.ShardID(), s.suid)
err = nil
}
}()
stat, err := s.Stats(ctx)
if err != nil && !errors.Is(err, apierr.ErrShardNoLeader) {
return err
}
if errors.Is(err, apierr.ErrShardNoLeader) || stat.LeaderDiskID != s.diskID {
return nil
}
units := stat.Units
if len(units) <= 1 {
return errors.New("no enough units to transfer leader")
}
for i := range units {
if units[i].GetDiskID() != s.diskID {
if err = s.TransferLeader(ctx, units[i].GetDiskID()); err != nil {
span.Errorf("shard[%d] suid[%d] transfer leader to disk[%d] failed: %s",
s.suid.ShardID(), s.suid, units[i].GetDiskID(), err.Error())
continue
}
span.Debugf("shard[%d] suid[%d] transfer leader to disk[%d] success",
s.suid.ShardID(), s.suid, units[i].GetDiskID())
return nil
}
}
return err
}
func (s *shard) SaveShardInfo(ctx context.Context, withLock bool, flush bool) error {
if withLock {
s.shardInfoMu.Lock()