mirror of
https://github.com/cubefs/cubefs.git
synced 2026-08-02 02:00:56 +00:00
feat(shardnode): update shard and remove shard implements
with #22357426 Signed-off-by: Cloudstriff <chenjiongwendao@qq.com>
This commit is contained in:
parent
38fd15c74c
commit
c06fdb8b3c
@ -2,9 +2,10 @@ package errors
|
||||
|
||||
// 2xx
|
||||
var (
|
||||
ErrShardNodeNotLeader = newError(1001, "shard node is not leader")
|
||||
ErrShardRangeMismatch = newError(1002, "shard range mismatch")
|
||||
ErrShardDoesNotExist = newError(1003, "shard doest not exist")
|
||||
ErrShardNodeDiskNotFound = newError(1004, "shard disk not found")
|
||||
ErrUnknownField = newError(1005, "unknown field")
|
||||
ErrShardNodeNotLeader = newError(1001, "shard node is not leader")
|
||||
ErrShardRangeMismatch = newError(1002, "shard range mismatch")
|
||||
ErrShardDoesNotExist = newError(1003, "shard doest not exist")
|
||||
ErrShardNodeDiskNotFound = newError(1004, "shard disk not found")
|
||||
ErrUnknownField = newError(1005, "unknown field")
|
||||
ErrShardRouteVersionNeedUpdate = newError(1006, "shard route version need update")
|
||||
)
|
||||
|
||||
@ -45,9 +45,9 @@ const (
|
||||
type ShardUpdateType uint8
|
||||
|
||||
const (
|
||||
ShardUpdateTypeAddMember ShardUpdateType = 0
|
||||
ShardUpdateTypeRemoveMember ShardUpdateType = 1
|
||||
ShardUpdateTypeSetNormal ShardUpdateType = 2
|
||||
ShardUpdateTypeAddMember ShardUpdateType = 1
|
||||
ShardUpdateTypeRemoveMember ShardUpdateType = 2
|
||||
ShardUpdateTypeUpdateMember ShardUpdateType = 3
|
||||
)
|
||||
|
||||
type FieldType uint8
|
||||
|
||||
@ -10,7 +10,6 @@ import (
|
||||
"github.com/cubefs/cubefs/blobstore/shardnode/base"
|
||||
"github.com/cubefs/cubefs/blobstore/shardnode/storage"
|
||||
"github.com/cubefs/cubefs/blobstore/util/closer"
|
||||
"github.com/cubefs/cubefs/blobstore/util/taskpool"
|
||||
)
|
||||
|
||||
const defaultTaskPoolSize = 64
|
||||
@ -30,7 +29,6 @@ type Catalog struct {
|
||||
routeVersion int64
|
||||
spaces sync.Map
|
||||
transport base.Transport
|
||||
taskPool taskpool.TaskPool
|
||||
|
||||
cfg *Config
|
||||
closer.Closer
|
||||
@ -42,7 +40,6 @@ func NewCatalog(ctx context.Context, cfg *Config) *Catalog {
|
||||
catalog := &Catalog{
|
||||
cfg: cfg,
|
||||
transport: cfg.Transport,
|
||||
taskPool: taskpool.New(defaultTaskPoolSize, defaultTaskPoolSize),
|
||||
Closer: closer.New(),
|
||||
}
|
||||
spaces, err := cfg.Transport.GetAllSpaces(ctx)
|
||||
|
||||
@ -107,7 +107,7 @@ func (s *service) loop(ctx context.Context) {
|
||||
continue
|
||||
}
|
||||
for _, task := range tasks {
|
||||
if err := s.catalog.ExecuteShardTask(ctx, task); err != nil {
|
||||
if err := s.executeShardTask(ctx, task); err != nil {
|
||||
span.Errorf("execute shard task[%+v] failed: %s", task, errors.Detail(err))
|
||||
continue
|
||||
}
|
||||
@ -140,3 +140,30 @@ func (s *service) getAlteredShardReports() []clustermgr.ShardReport {
|
||||
|
||||
return ret
|
||||
}
|
||||
|
||||
func (s *service) executeShardTask(ctx context.Context, task clustermgr.ShardTask) error {
|
||||
span := trace.SpanFromContext(ctx)
|
||||
|
||||
disk, err := s.getDisk(task.DiskID)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
shard, err := disk.GetShard(task.Suid)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
switch task.TaskType {
|
||||
case proto.ShardTaskTypeClearShard:
|
||||
s.taskPool.Run(func() {
|
||||
if shard.GetEpoch() == task.Epoch {
|
||||
err := disk.DeleteShard(ctx, task.Suid)
|
||||
if err != nil {
|
||||
span.Errorf("delete shard task[%+v] failed: %s", task, err)
|
||||
}
|
||||
}
|
||||
})
|
||||
default:
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
@ -225,7 +225,6 @@ func (d *Disk) AddShard(ctx context.Context, suid proto.Suid,
|
||||
return nil
|
||||
}
|
||||
|
||||
// todo: update shard Units support
|
||||
func (d *Disk) UpdateShard(ctx context.Context, suid proto.Suid, op proto.ShardUpdateType, node clustermgr.ShardUnitInfo) error {
|
||||
shard, err := d.GetShard(suid)
|
||||
if err != nil {
|
||||
@ -237,22 +236,7 @@ func (d *Disk) UpdateShard(ctx context.Context, suid proto.Suid, op proto.ShardU
|
||||
return err
|
||||
}
|
||||
|
||||
switch op {
|
||||
case proto.ShardUpdateTypeAddMember:
|
||||
shard.raftGroup.MemberChange(ctx, &raft.Member{
|
||||
NodeID: uint64(node.DiskID),
|
||||
Host: nodeHost.String(),
|
||||
Type: raft.MemberChangeType_AddMember,
|
||||
Learner: node.Learner,
|
||||
})
|
||||
case proto.ShardUpdateTypeRemoveMember:
|
||||
shard.raftGroup.MemberChange(ctx, &raft.Member{
|
||||
NodeID: uint64(node.DiskID),
|
||||
Host: nodeHost.String(),
|
||||
Type: raft.MemberChangeType_RemoveMember,
|
||||
Learner: node.Learner,
|
||||
})
|
||||
}
|
||||
shard.UpdateShard(ctx, op, node, nodeHost.String())
|
||||
|
||||
return nil
|
||||
}
|
||||
@ -269,17 +253,37 @@ func (d *Disk) GetShard(suid proto.Suid) (*shard, error) {
|
||||
}
|
||||
|
||||
func (d *Disk) DeleteShard(ctx context.Context, suid proto.Suid) error {
|
||||
d.shardsMu.Lock()
|
||||
d.shardsMu.RLock()
|
||||
shard := d.shardsMu.shards[suid]
|
||||
delete(d.shardsMu.shards, suid)
|
||||
d.shardsMu.Unlock()
|
||||
d.shardsMu.RUnlock()
|
||||
|
||||
if shard != nil {
|
||||
shard.Stop()
|
||||
shard.Close()
|
||||
// todo: clear shard's data
|
||||
if shard == nil {
|
||||
return nil
|
||||
}
|
||||
|
||||
d.shardsMu.Lock()
|
||||
defer d.shardsMu.Unlock()
|
||||
|
||||
if shard == nil {
|
||||
return nil
|
||||
}
|
||||
|
||||
nodeHost, err := d.cfg.RaftConfig.Resolver.Resolve(ctx, uint64(d.diskInfo.DiskID))
|
||||
if err != nil {
|
||||
return errors.Info(err, "resolve disk node host failed")
|
||||
}
|
||||
|
||||
if err := shard.DeleteShard(ctx, nodeHost.String()); err != nil {
|
||||
return errors.Info(err, "delete shard failed")
|
||||
}
|
||||
|
||||
// remove raft group
|
||||
if err := d.raftManager.RemoveRaftGroup(ctx, uint64(suid.ShardID()), true); err != nil {
|
||||
return errors.Info(err, "remove raft group failed")
|
||||
}
|
||||
|
||||
delete(d.shardsMu.shards, suid)
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
|
||||
@ -126,7 +126,10 @@ func TestServerDisk_Shard(t *testing.T) {
|
||||
_, err = d.d.GetShard(suid2)
|
||||
require.NoError(t, err)
|
||||
|
||||
d.d.RangeShard(func(s ShardHandler) bool { s.Checkpoint(ctx); return true })
|
||||
d.d.RangeShard(func(s ShardHandler) bool {
|
||||
require.NoError(t, s.Checkpoint(ctx))
|
||||
return true
|
||||
})
|
||||
require.NoError(t, d.d.DeleteShard(ctx, suid2))
|
||||
require.NoError(t, d.d.DeleteShard(ctx, suid2))
|
||||
|
||||
|
||||
@ -153,7 +153,7 @@ func (a *addressResolver) Resolve(ctx context.Context, diskID uint64) (raft.Addr
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return nodeAddr{addr: node.Host}, nil
|
||||
return nodeAddr{addr: node.RaftHost}, nil
|
||||
}
|
||||
|
||||
type nodeAddr struct {
|
||||
|
||||
@ -3,9 +3,11 @@ package storage
|
||||
import (
|
||||
"context"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"golang.org/x/sync/singleflight"
|
||||
|
||||
"github.com/cubefs/cubefs/blobstore/api/clustermgr"
|
||||
"github.com/cubefs/cubefs/blobstore/api/shardnode"
|
||||
apierr "github.com/cubefs/cubefs/blobstore/common/errors"
|
||||
kvstore "github.com/cubefs/cubefs/blobstore/common/kvstorev2"
|
||||
@ -20,6 +22,11 @@ import (
|
||||
|
||||
const keyLocksNum = 1024
|
||||
|
||||
const (
|
||||
shardStatusNormal = shardStatus(1)
|
||||
shardStatusStopReadWrite = shardStatus(2)
|
||||
)
|
||||
|
||||
type (
|
||||
ShardHandler interface {
|
||||
InsertItem(ctx context.Context, h OpHeader, i shardnode.Item) error
|
||||
@ -42,6 +49,12 @@ type (
|
||||
TruncateWalLogInterval uint64 `json:"truncate_wal_log_interval"`
|
||||
}
|
||||
|
||||
ShardStats struct {
|
||||
Leader proto.DiskID
|
||||
Epoch uint64
|
||||
Units []shardUnitInfo
|
||||
}
|
||||
|
||||
shardConfig struct {
|
||||
*ShardBaseConfig
|
||||
suid proto.Suid
|
||||
@ -52,11 +65,7 @@ type (
|
||||
addrResolver raft.AddressResolver
|
||||
}
|
||||
|
||||
ShardStats struct {
|
||||
Leader proto.DiskID
|
||||
Epoch uint64
|
||||
Units []shardUnitInfo
|
||||
}
|
||||
shardStatus uint8
|
||||
)
|
||||
|
||||
func newShard(ctx context.Context, cfg shardConfig) (s *shard, err error) {
|
||||
@ -67,15 +76,13 @@ func newShard(ctx context.Context, cfg shardConfig) (s *shard, err error) {
|
||||
suid: cfg.suid,
|
||||
diskID: cfg.diskID,
|
||||
|
||||
// startIno: calculateStartIno(cfg.shardInfo.ShardID),
|
||||
|
||||
store: cfg.store,
|
||||
shardKeys: &shardKeysGenerator{
|
||||
suid: cfg.suid,
|
||||
},
|
||||
cfg: cfg.ShardBaseConfig,
|
||||
}
|
||||
s.shardMu.shardInfo = cfg.shardInfo
|
||||
s.shardInfoMu.shardInfo = cfg.shardInfo
|
||||
|
||||
learner := false
|
||||
for _, node := range cfg.shardInfo.Units {
|
||||
@ -101,6 +108,7 @@ func newShard(ctx context.Context, cfg shardConfig) (s *shard, err error) {
|
||||
span.Debugf("shard members: %+v", members)
|
||||
|
||||
s.raftGroup, err = cfg.raftManager.CreateRaftGroup(context.Background(), &raft.GroupConfig{
|
||||
// Note: set raft group id with shard id as all shard node share the same shard id
|
||||
ID: uint64(cfg.shardInfo.ShardID),
|
||||
Applied: cfg.shardInfo.AppliedIndex,
|
||||
Members: members,
|
||||
@ -118,18 +126,20 @@ func newShard(ctx context.Context, cfg shardConfig) (s *shard, err error) {
|
||||
}
|
||||
|
||||
type shard struct {
|
||||
suid proto.Suid
|
||||
diskID proto.DiskID
|
||||
lastStableIndex uint64
|
||||
lastTruncatedIndex uint64
|
||||
suid proto.Suid
|
||||
diskID proto.DiskID
|
||||
|
||||
shardMu struct {
|
||||
shardState shardState
|
||||
shardInfoMu struct {
|
||||
sync.RWMutex
|
||||
shardInfo
|
||||
leader proto.DiskID
|
||||
}
|
||||
sf singleflight.Group
|
||||
|
||||
leader proto.DiskID
|
||||
lastStableIndex uint64
|
||||
lastTruncatedIndex uint64
|
||||
}
|
||||
|
||||
sf singleflight.Group
|
||||
shardKeys *shardKeysGenerator
|
||||
store *store.Store
|
||||
raftGroup raft.Group
|
||||
@ -144,7 +154,7 @@ func (s *shard) InsertItem(ctx context.Context, h OpHeader, i shardnode.Item) er
|
||||
return err
|
||||
}
|
||||
|
||||
internalItem := s.protoItemToInternalItem(i)
|
||||
internalItem := protoItemToInternalItem(i)
|
||||
data, err := internalItem.Marshal()
|
||||
if err != nil {
|
||||
return err
|
||||
@ -168,7 +178,7 @@ func (s *shard) UpdateItem(ctx context.Context, h OpHeader, i shardnode.Item) er
|
||||
return err
|
||||
}
|
||||
|
||||
internalItem := s.protoItemToInternalItem(i)
|
||||
internalItem := protoItemToInternalItem(i)
|
||||
data, err := internalItem.Marshal()
|
||||
if err != nil {
|
||||
return err
|
||||
@ -219,6 +229,30 @@ func (s *shard) GetItem(ctx context.Context, h OpHeader, id []byte) (protoItem s
|
||||
return
|
||||
}
|
||||
|
||||
func (s *shard) GetItems(ctx context.Context, keys [][]byte) (ret []shardnode.Item, err error) {
|
||||
store := s.store.KVStore()
|
||||
vgs, err := store.MultiGet(ctx, dataCF, keys, nil)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
|
||||
ret = make([]shardnode.Item, len(vgs))
|
||||
for i := range ret {
|
||||
item := &item{}
|
||||
err = item.Unmarshal(vgs[i].Value())
|
||||
vgs[i].Close()
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
ret[i] = shardnode.Item{
|
||||
ID: item.ID,
|
||||
Fields: internalFieldsToProtoFields(item.Fields),
|
||||
}
|
||||
}
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
func (s *shard) ListItem(ctx context.Context, h OpHeader, prefix, marker []byte, count uint64) (items []shardnode.Item, nextMarker []byte, err error) {
|
||||
if err := s.checkShardOptHeader(h); err != nil {
|
||||
return nil, nil, err
|
||||
@ -261,19 +295,19 @@ func (s *shard) ListItem(ctx context.Context, h OpHeader, prefix, marker []byte,
|
||||
}
|
||||
|
||||
func (s *shard) GetEpoch() uint64 {
|
||||
s.shardMu.RLock()
|
||||
epoch := s.shardMu.Epoch
|
||||
s.shardMu.RUnlock()
|
||||
s.shardInfoMu.RLock()
|
||||
epoch := s.shardInfoMu.Epoch
|
||||
s.shardInfoMu.RUnlock()
|
||||
return epoch
|
||||
}
|
||||
|
||||
func (s *shard) Stats() ShardStats {
|
||||
s.shardMu.RLock()
|
||||
replicates := make([]shardUnitInfo, len(s.shardMu.Units))
|
||||
copy(replicates, s.shardMu.Units)
|
||||
epoch := s.shardMu.Epoch
|
||||
leader := s.shardMu.leader
|
||||
s.shardMu.RUnlock()
|
||||
s.shardInfoMu.RLock()
|
||||
replicates := make([]shardUnitInfo, len(s.shardInfoMu.Units))
|
||||
copy(replicates, s.shardInfoMu.Units)
|
||||
epoch := s.shardInfoMu.Epoch
|
||||
leader := s.shardInfoMu.leader
|
||||
s.shardInfoMu.RUnlock()
|
||||
|
||||
return ShardStats{
|
||||
Leader: leader,
|
||||
@ -282,19 +316,6 @@ func (s *shard) Stats() ShardStats {
|
||||
}
|
||||
}
|
||||
|
||||
func (s *shard) Start() {
|
||||
}
|
||||
|
||||
func (s *shard) Stop() {
|
||||
// TODO: stop all operation on this shard
|
||||
}
|
||||
|
||||
func (s *shard) Close() {
|
||||
// TODO: wait all operation done on this shard and then close shard, ensure memory safe
|
||||
|
||||
s.raftGroup.Close()
|
||||
}
|
||||
|
||||
// Checkpoint do checkpoint job with raft group
|
||||
// we should do any memory flush job or dump worker here
|
||||
func (s *shard) Checkpoint(ctx context.Context) error {
|
||||
@ -306,28 +327,47 @@ func (s *shard) Checkpoint(ctx context.Context) error {
|
||||
}
|
||||
|
||||
// truncate raft log finally
|
||||
if appliedIndex > s.lastTruncatedIndex+s.cfg.TruncateWalLogInterval*2 {
|
||||
if appliedIndex > s.shardInfoMu.lastTruncatedIndex+s.cfg.TruncateWalLogInterval*2 {
|
||||
if err := s.raftGroup.Truncate(ctx, appliedIndex-s.cfg.TruncateWalLogInterval); err != nil {
|
||||
return errors.Info(err, "truncate raft wal log failed")
|
||||
}
|
||||
s.lastTruncatedIndex = appliedIndex - s.cfg.TruncateWalLogInterval
|
||||
s.shardInfoMu.lastTruncatedIndex = appliedIndex - s.cfg.TruncateWalLogInterval
|
||||
}
|
||||
|
||||
// save last stable index
|
||||
s.lastStableIndex = appliedIndex
|
||||
s.shardInfoMu.lastStableIndex = appliedIndex
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *shard) UpdateShard(ctx context.Context, op proto.ShardUpdateType, node clustermgr.ShardUnitInfo, nodeHost string) {
|
||||
switch op {
|
||||
case proto.ShardUpdateTypeAddMember, proto.ShardUpdateTypeUpdateMember:
|
||||
s.raftGroup.MemberChange(ctx, &raft.Member{
|
||||
NodeID: uint64(node.DiskID),
|
||||
Host: nodeHost,
|
||||
Type: raft.MemberChangeType_AddMember,
|
||||
Learner: node.Learner,
|
||||
})
|
||||
case proto.ShardUpdateTypeRemoveMember:
|
||||
s.raftGroup.MemberChange(ctx, &raft.Member{
|
||||
NodeID: uint64(node.DiskID),
|
||||
Host: nodeHost,
|
||||
Type: raft.MemberChangeType_RemoveMember,
|
||||
Learner: node.Learner,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func (s *shard) SaveShardInfo(ctx context.Context, withLock bool, flush bool) error {
|
||||
if withLock {
|
||||
s.shardMu.Lock()
|
||||
defer s.shardMu.Unlock()
|
||||
s.shardInfoMu.Lock()
|
||||
defer s.shardInfoMu.Unlock()
|
||||
}
|
||||
|
||||
kvStore := s.store.KVStore()
|
||||
key := s.shardKeys.encodeShardInfoKey()
|
||||
value, err := s.shardMu.shardInfo.Marshal()
|
||||
value, err := s.shardInfoMu.shardInfo.Marshal()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
@ -346,63 +386,207 @@ func (s *shard) SaveShardInfo(ctx context.Context, withLock bool, flush bool) er
|
||||
return kvStore.FlushCF(ctx, dataCF)
|
||||
}
|
||||
|
||||
func (s *shard) DeleteShard(ctx context.Context, nodeHost string) error {
|
||||
// 1. check raft status and remove member itself
|
||||
stat, err := s.raftGroup.Stat()
|
||||
if err != nil {
|
||||
return errors.Info(err, "raft stat failed")
|
||||
}
|
||||
if len(stat.Nodes) > 1 {
|
||||
for _, diskID := range stat.Nodes {
|
||||
if uint64(s.diskID) == diskID {
|
||||
if err := s.raftGroup.MemberChange(ctx, &raft.Member{
|
||||
NodeID: uint64(s.diskID),
|
||||
Host: nodeHost,
|
||||
Type: raft.MemberChangeType_RemoveMember,
|
||||
}); err != nil {
|
||||
return errors.Info(err, "remove raft member failed")
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// 2. stop all writing in this shard
|
||||
if err := s.Stop(); err != nil {
|
||||
return errors.Info(err, "stop shard failed")
|
||||
}
|
||||
|
||||
// 3. clear all shard's data
|
||||
kvStore := s.store.KVStore()
|
||||
batch := kvStore.NewWriteBatch()
|
||||
|
||||
batch.DeleteRange(dataCF, s.shardKeys.encodeShardDataPrefix(), s.shardKeys.encodeShardDataMaxPrefix())
|
||||
batch.Delete(dataCF, s.shardKeys.encodeShardInfoKey())
|
||||
if err := kvStore.Write(ctx, batch); err != nil {
|
||||
return errors.Info(err, "kvstore write batch failed")
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *shard) Start() {
|
||||
}
|
||||
|
||||
func (s *shard) Stop() error {
|
||||
if err := s.shardState.stopWriting(); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
// wait all operation done on this shard before close shard, ensure memory safe
|
||||
s.shardState.waitPendingRequestDone()
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *shard) Close() {
|
||||
s.raftGroup.Close()
|
||||
}
|
||||
|
||||
func (s *shard) GetAppliedIndex() uint64 {
|
||||
return (*shardSM)(s).getAppliedIndex()
|
||||
}
|
||||
|
||||
func (s *shard) GetStableIndex() uint64 {
|
||||
return s.lastStableIndex
|
||||
}
|
||||
|
||||
func (s *shard) GetItems(ctx context.Context, keys [][]byte) (ret []shardnode.Item, err error) {
|
||||
store := s.store.KVStore()
|
||||
vgs, err := store.MultiGet(ctx, dataCF, keys, nil)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
|
||||
ret = make([]shardnode.Item, len(vgs))
|
||||
for i := range ret {
|
||||
item := &item{}
|
||||
err = item.Unmarshal(vgs[i].Value())
|
||||
vgs[i].Close()
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
ret[i] = shardnode.Item{
|
||||
ID: item.ID,
|
||||
Fields: internalFieldsToProtoFields(item.Fields),
|
||||
}
|
||||
}
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
func (s *shard) protoItemToInternalItem(i shardnode.Item) (ret item) {
|
||||
ret = item{
|
||||
ID: i.ID,
|
||||
Fields: protoFieldsToInternalFields(i.Fields),
|
||||
}
|
||||
return
|
||||
return s.shardInfoMu.lastStableIndex
|
||||
}
|
||||
|
||||
func (s *shard) checkShardOptHeader(h OpHeader) error {
|
||||
// todo: check shard route version ?
|
||||
ci := sharding.NewCompareItem(s.shardMu.Range.Type, h.ShardKeys)
|
||||
if !s.shardMu.Range.Belong(ci) {
|
||||
ci := sharding.NewCompareItem(s.shardInfoMu.Range.Type, h.ShardKeys)
|
||||
if !s.shardInfoMu.Range.Belong(ci) {
|
||||
return apierr.ErrShardRangeMismatch
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *shard) isLeader() bool {
|
||||
s.shardMu.RLock()
|
||||
isLeader := s.shardMu.leader == s.diskID
|
||||
s.shardMu.RUnlock()
|
||||
s.shardInfoMu.RLock()
|
||||
isLeader := s.shardInfoMu.leader == s.diskID
|
||||
s.shardInfoMu.RUnlock()
|
||||
|
||||
return isLeader
|
||||
}
|
||||
|
||||
type shardState struct {
|
||||
status shardStatus
|
||||
pendingReqs int
|
||||
pendingReqsWg sync.WaitGroup
|
||||
|
||||
splitting bool
|
||||
lastSplitTime time.Time
|
||||
splitDone chan struct{}
|
||||
|
||||
lock sync.RWMutex
|
||||
}
|
||||
|
||||
func (s *shardState) stopWriting() error {
|
||||
s.lock.Lock()
|
||||
defer s.lock.Unlock()
|
||||
|
||||
// shard is splitting, can't be stop writing
|
||||
if s.splitting {
|
||||
return errors.New("shard is splitting")
|
||||
}
|
||||
|
||||
if s.status != shardStatusStopReadWrite {
|
||||
s.pendingReqsWg.Add(s.pendingReqs)
|
||||
s.status = shardStatusStopReadWrite
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *shardState) splitStartWriting() {
|
||||
s.lock.Lock()
|
||||
s.status = shardStatusNormal
|
||||
s.lock.Unlock()
|
||||
|
||||
close(s.splitDone)
|
||||
}
|
||||
|
||||
func (s *shardState) splitStopWriting() {
|
||||
s.lock.Lock()
|
||||
defer s.lock.Unlock()
|
||||
|
||||
if s.status != shardStatusStopReadWrite {
|
||||
s.pendingReqsWg.Add(s.pendingReqs)
|
||||
s.status = shardStatusStopReadWrite
|
||||
}
|
||||
|
||||
s.splitDone = make(chan struct{})
|
||||
}
|
||||
|
||||
func (s *shardState) prepRWCheck() error {
|
||||
s.lock.Lock()
|
||||
|
||||
// allow writing check in the list lock arena
|
||||
if !s.allowRW() {
|
||||
s.lock.Unlock()
|
||||
|
||||
if !s.splitting {
|
||||
// return route need update when unbind shard stop write progress
|
||||
return apierr.ErrShardRouteVersionNeedUpdate
|
||||
}
|
||||
|
||||
// wait for start write when split shard progress
|
||||
s.waitSplitDone()
|
||||
s.lock.Lock()
|
||||
}
|
||||
|
||||
s.pendingReqs++
|
||||
s.lock.Unlock()
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *shardState) prepRWCheckDone() {
|
||||
s.lock.Lock()
|
||||
s.pendingReqs--
|
||||
// decrease pending request
|
||||
if !s.allowRW() {
|
||||
s.pendingReqsWg.Done()
|
||||
}
|
||||
s.lock.Unlock()
|
||||
}
|
||||
|
||||
func (s *shardState) waitSplitDone() {
|
||||
s.lock.RLock()
|
||||
done := s.splitDone
|
||||
s.lock.RUnlock()
|
||||
<-done
|
||||
}
|
||||
|
||||
func (s *shardState) startSplitting() bool {
|
||||
s.lock.Lock()
|
||||
defer s.lock.Unlock()
|
||||
|
||||
// already in splitting, return false and do not start new split
|
||||
if s.splitting {
|
||||
return false
|
||||
}
|
||||
// check if shard stop writing first
|
||||
if !s.allowRW() {
|
||||
return false
|
||||
}
|
||||
|
||||
s.splitting = true
|
||||
return true
|
||||
}
|
||||
|
||||
func (s *shardState) stopSplitting() {
|
||||
s.lock.Lock()
|
||||
s.lastSplitTime = time.Now()
|
||||
s.splitting = false
|
||||
s.lock.Unlock()
|
||||
}
|
||||
|
||||
func (s *shardState) waitPendingRequestDone() {
|
||||
s.pendingReqsWg.Wait()
|
||||
}
|
||||
|
||||
func (s *shardState) allowRW() bool {
|
||||
return s.status != shardStatusStopReadWrite
|
||||
}
|
||||
|
||||
type shardKeysGenerator struct {
|
||||
suid proto.Suid
|
||||
}
|
||||
@ -447,6 +631,19 @@ func (s *shardKeysGenerator) encodeShardDataMaxPrefix() []byte {
|
||||
return key
|
||||
}
|
||||
|
||||
type shardStopper struct {
|
||||
stopWriteDone chan struct{}
|
||||
pendingReqWg sync.WaitGroup
|
||||
}
|
||||
|
||||
func protoItemToInternalItem(i shardnode.Item) (ret item) {
|
||||
ret = item{
|
||||
ID: i.ID,
|
||||
Fields: protoFieldsToInternalFields(i.Fields),
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
func protoFieldsToInternalFields(external []shardnode.Field) []storageproto.Field {
|
||||
// todo: use memory pool
|
||||
ret := make([]storageproto.Field, len(external))
|
||||
|
||||
@ -60,9 +60,9 @@ func (s *shardSM) Apply(cxt context.Context, pd []raft.ProposalData, index uint6
|
||||
func (s *shardSM) LeaderChange(peerID uint64) error {
|
||||
log.Info("shard receive Leader change", peerID)
|
||||
// todo: report Leader change to master
|
||||
s.shardMu.Lock()
|
||||
s.shardMu.leader = proto.DiskID(peerID)
|
||||
s.shardMu.Unlock()
|
||||
s.shardInfoMu.Lock()
|
||||
s.shardInfoMu.leader = proto.DiskID(peerID)
|
||||
s.shardInfoMu.Unlock()
|
||||
// todo: read index before start to serve request
|
||||
|
||||
return nil
|
||||
@ -71,28 +71,29 @@ func (s *shardSM) LeaderChange(peerID uint64) error {
|
||||
func (s *shardSM) ApplyMemberChange(cc *raft.Member, index uint64) error {
|
||||
_, c := trace.StartSpanFromContext(context.Background(), "")
|
||||
|
||||
s.shardMu.Lock()
|
||||
defer s.shardMu.Unlock()
|
||||
s.shardInfoMu.Lock()
|
||||
defer s.shardInfoMu.Unlock()
|
||||
|
||||
switch cc.Type {
|
||||
case raft.MemberChangeType_AddMember:
|
||||
found := false
|
||||
for _, node := range s.shardMu.Units {
|
||||
if node.DiskID == proto.DiskID(cc.NodeID) {
|
||||
for i := range s.shardInfoMu.Units {
|
||||
if s.shardInfoMu.Units[i].DiskID == proto.DiskID(cc.NodeID) {
|
||||
s.shardInfoMu.Units[i].Learner = cc.Learner
|
||||
found = true
|
||||
break
|
||||
}
|
||||
}
|
||||
if !found {
|
||||
s.shardMu.Units = append(s.shardMu.Units, clustermgr.ShardUnitInfo{
|
||||
s.shardInfoMu.Units = append(s.shardInfoMu.Units, clustermgr.ShardUnitInfo{
|
||||
DiskID: proto.DiskID(cc.NodeID),
|
||||
Learner: cc.Learner,
|
||||
})
|
||||
}
|
||||
case raft.MemberChangeType_RemoveMember:
|
||||
for i, node := range s.shardMu.Units {
|
||||
for i, node := range s.shardInfoMu.Units {
|
||||
if node.DiskID == proto.DiskID(cc.NodeID) {
|
||||
s.shardMu.Units = append(s.shardMu.Units[:i], s.shardMu.Units[i+1:]...)
|
||||
s.shardInfoMu.Units = append(s.shardInfoMu.Units[:i], s.shardInfoMu.Units[i+1:]...)
|
||||
break
|
||||
}
|
||||
}
|
||||
@ -268,9 +269,9 @@ func (s *shardSM) applyDeleteItem(ctx context.Context, data []byte) error {
|
||||
}
|
||||
|
||||
func (s *shardSM) setAppliedIndex(index uint64) {
|
||||
atomic.StoreUint64(&s.shardMu.AppliedIndex, index)
|
||||
atomic.StoreUint64(&s.shardInfoMu.AppliedIndex, index)
|
||||
}
|
||||
|
||||
func (s *shardSM) getAppliedIndex() uint64 {
|
||||
return atomic.LoadUint64(&s.shardMu.AppliedIndex)
|
||||
return atomic.LoadUint64(&s.shardInfoMu.AppliedIndex)
|
||||
}
|
||||
|
||||
@ -71,10 +71,13 @@ func newMockShard(tb testing.TB) (*mockShard, func()) {
|
||||
BatchInflightSize: 1 << 20,
|
||||
},
|
||||
},
|
||||
shardMu: struct {
|
||||
shardInfoMu: struct {
|
||||
sync.RWMutex
|
||||
shardInfo
|
||||
leader proto.DiskID
|
||||
|
||||
leader proto.DiskID
|
||||
lastStableIndex uint64
|
||||
lastTruncatedIndex uint64
|
||||
}{
|
||||
leader: 1, shardInfo: shardInfo{
|
||||
ShardID: 1,
|
||||
|
||||
@ -3,6 +3,8 @@ package shardnode
|
||||
import (
|
||||
"sync"
|
||||
|
||||
"github.com/cubefs/cubefs/blobstore/util/taskpool"
|
||||
|
||||
"github.com/cubefs/cubefs/blobstore/shardnode/base"
|
||||
|
||||
apierr "github.com/cubefs/cubefs/blobstore/common/errors"
|
||||
@ -33,6 +35,7 @@ type service struct {
|
||||
catalog *catalog.Catalog
|
||||
disks map[proto.DiskID]*storage.Disk
|
||||
transport base.Transport
|
||||
taskPool taskpool.TaskPool
|
||||
|
||||
cfg Config
|
||||
lock sync.RWMutex
|
||||
|
||||
Loading…
Reference in New Issue
Block a user