在这部分中我们要在lab4的基础上更近一步,lab4是只有一组raft服务,lab5就是Multiraft(multi-group),进一步加入了分片的机制,主要用于解决在大规模分布式系统中高可用、强一致性、可扩展性等问题。

单 Raft Group 的局限性:

Raft 是一种用于实现分布式系统一致性的共识算法,通常以一个“Raft Group”(即一组节点共同维护一个日志副本)为单位工作。但在实际应用中,如果整个系统只用一个 Raft Group,有以下缺点:

  • 写性能瓶颈: 所有的写请求必须由 Leader 单节点处理并定序,无法利用集群多节点的并发能力,吞吐量受限于单机资源。
  • 扩展性受限: 增加节点仅能提升容灾能力,无法线性提升写性能,且集群总存储容量被限制在单机磁盘大小。
  • 故障恢复慢与热点问题: 巨大的状态机导致快照生成和日志回放缓慢,故障恢复时间长;且单一 Group 难以进行细粒度的负载均衡,易形成读写热点。

分片 + 多 Raft Group:

Multiraft 将整个数据集划分为多个分片(shard),每个分片由一个独立的 Raft Group 负责管理。

  • 每个 Raft Group 只负责一部分数据,互不影响;
  • 写/读请求可以并行处理,提高系统吞吐;
  • 系统可以通过增加分片数量来横向扩展;
  • 故障隔离:某个 Group 出现问题(如网络分区、节点宕机),不会影响其他 Group。

shardctrler

针对5A的实验,主要就是设计一个shardctrler,它的作用是维护集群的配置信息(Configuration),即负责管理“分片(Shard)”到“复制组(Replica Group / GID)”的映射关系,并在集群拓扑变化时自动进行分片的负载均衡(Rebalancing)。

具体来说,它需要通过 Raft 保证高可用和强一致性,并支持以下 4 个核心操作(RPC):

  1. Join (加入新组): 当新的复制组(Replica Group)加入集群时,Ctrler 需要将部分分片从现有组迁移给新组,以实现负载均衡。
  2. Leave (移出旧组): 当某个复制组想要离开集群时,Ctrler 需要将其持有的分片重新分配给剩余的组,确保数据不丢失且分布均匀。
  3. Move (手动迁移): 允许管理员强制将某个特定的分片(Shard)分配给特定的组(GID),主要用于调试或处理热点。
  4. Query (查询配置): 允许 ShardKV 节点或 Client 查询最新的或历史版本的配置(Config),以便知道读写请求该发往哪里。

它本身一个group,所以基本的接收命令然后提交等操作和之前完全一样,主要就是这四个操作和rebalance,注意rebalance要最少移动原则。

  1
  2
  3
  4
  5
  6
  7
  8
  9
 10
 11
 12
 13
 14
 15
 16
 17
 18
 19
 20
 21
 22
 23
 24
 25
 26
 27
 28
 29
 30
 31
 32
 33
 34
 35
 36
 37
 38
 39
 40
 41
 42
 43
 44
 45
 46
 47
 48
 49
 50
 51
 52
 53
 54
 55
 56
 57
 58
 59
 60
 61
 62
 63
 64
 65
 66
 67
 68
 69
 70
 71
 72
 73
 74
 75
 76
 77
 78
 79
 80
 81
 82
 83
 84
 85
 86
 87
 88
 89
 90
 91
 92
 93
 94
 95
 96
 97
 98
 99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
// 计算平均值,先剔除多余的或无效的分片,再将空闲分片分配给需要的新组
func (sc *ShardCtrler) rebalance(config *Config) {
	allNum := len(config.Groups)

	if allNum == 0 {
		for i := 0; i < NShards; i++ {
			config.Shards[i] = 0
		}
		return
	}

	var keys []int
	for key := range config.Groups {
		keys = append(keys, key)
	}
	sort.Ints(keys)

	avg := max((NShards / allNum), 1)
	remainder := NShards % allNum

	cnt := make(map[int]int, allNum)
	for i, gid := range keys {
		if i < remainder {
			cnt[gid] = avg + 1
		} else {
			cnt[gid] = avg
		}
	}

	for i := 0; i < NShards; i++ {
		gid := config.Shards[i]
		if gid == 0 {
			continue
		}
		if _, ok := config.Groups[gid]; !ok {
			config.Shards[i] = 0
			continue
		}
		if cnt[gid] == 0 {
			config.Shards[i] = 0
		} else {
			cnt[gid] -= 1
		}
	}

	freeG := list.New()

	type remain struct {
		gid int
		cnt int
	}

	for _, k := range keys {
		if cnt[k] > 0 {
			freeG.PushBack(remain{gid: k, cnt: cnt[k]})
		}
	}

	for i := 0; i < NShards; i++ {
		if config.Shards[i] == 0 {
			front := freeG.Front()
			value := front.Value.(remain)
			config.Shards[i] = value.gid
			if value.cnt == 1 {
				freeG.Remove(front)
			} else {
				value.cnt -= 1
				front.Value = value
			}
		}
	}

}

func (sc *ShardCtrler) getNewConfig() *Config {
	n := len(sc.configs)
	config := sc.configs[n-1]
	newConfig := config
	newConfig.Groups = make(map[int][]string, len(config.Groups))

	for k, v := range config.Groups {
		newConfig.Groups[k] = v
	}
	newConfig.Num = config.Num + 1
	return &newConfig
}

func (sc *ShardCtrler) applyJoin(servers map[int][]string) {
	newConfig := sc.getNewConfig()

	for k, v := range servers {
		newConfig.Groups[k] = v
	}

	sc.rebalance(newConfig)
	sc.configs = append(sc.configs, *newConfig)
}

func (sc *ShardCtrler) applyLeave(LeaveGIDs []int) {
	newConfig := sc.getNewConfig()

	for _, gid := range LeaveGIDs {
		if _, ok := newConfig.Groups[gid]; ok {
			delete(newConfig.Groups, gid)
		}
	}

	sc.rebalance(newConfig)
	sc.configs = append(sc.configs, *newConfig)
}

func (sc *ShardCtrler) applyMove(gid int, shard int) {
	newConfig := sc.getNewConfig()
	newConfig.Shards[shard] = gid
	sc.configs = append(sc.configs, *newConfig)
}

func (sc *ShardCtrler) applyQuery(num int) *Config {
	lastConfigIndex := len(sc.configs) - 1

	if num == -1 || num > lastConfigIndex {
		return &sc.configs[lastConfigIndex]
	} else {
		return &sc.configs[num]
	}
}

Shard Movement

上一部分的实验(5A)主要完成了控制面的工作,即通过 ShardCtrler 维护集群配置和分片映射,但并没有涉及真正的数据面迁移。 这一部分(5B)的核心就是分片移动(Shard Movement):当配置发生变更时,节点之间需要通过网络传输分片数据。基础的 KV 存储逻辑可以复用 kvraft 的代码,但我们需要在其之上构建一套复杂的状态机流转机制

分片状态

由于每一组group负责的分片不同了,因此我们定义一下分片,来判断server是否对该分片提供服务,分片的状态定义如下:

1
2
3
4
5
6
const (
	Serving   ShardStatus = iota // 正常服务中
	Pulling                      // 正在从别人那拉数据 (不可读写)
	BePulling                    // 正在被别人拉数据
	GCing                        // 等待清理
)
 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
func (kv *ShardKV) checkShard(key string) Err {
	shardID := key2shard(key)
	kv.mu.Lock()
	defer kv.mu.Unlock()

	if kv.currentConfig.Shards[shardID] != kv.gid {
		return ErrWrongGroup
	}

	if shard, ok := kv.data[shardID]; !ok || shard.State != Serving {
		return ErrNotServing
	}
	DPrintf(dAsk, "%v: have a ask, shardID is %d, pass check\n", kv, shardID)
	return OK
}

每一个分片到来前我们都要进行如上的检查,如果shard对应的owner不对或者不是serving状态下,就返回错误。这里其实还可以优化,就是旧owner如果在BePulling,完全可以提供读服务,但写一定要拒绝。

Config Update

接下来就考虑分片的迁移:每个 Replica Group 都需要感知全局配置的变化,因此启动一个后台协程(Goroutine)定期监控配置变更(Config Polling)是整个迁移流程的起点。

后台监控线程不能直接把配置从 Config 1 跳到 Config 10。必须逐号更新。因为我们需要知道分片是从哪里来的。如果是从1跳到10,中间可能经历了多次变更,我们将丢失分片的“前任持有者”信息,导致无法去拉取数据。

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
func (kv *ShardKV) monitorConfig() {
	for !kv.killed() {
		if kv.isLeader() {
			kv.mu.Lock()
			currentNum := kv.currentConfig.Num
			if !kv.allShardsReady() {
				kv.mu.Unlock()
				time.Sleep(100 * time.Millisecond)
				continue
			}
			kv.mu.Unlock()

			newConfig := kv.ctrClients.Query(currentNum + 1)
			DPrintf(dConfig, "%v: all shard ready, query next Config, next num is %d\n", kv, newConfig.Num)

			if newConfig.Num == currentNum+1 {
				kv.mu.Lock()
				if kv.currentConfig.Num == currentNum {
					kv.mu.Unlock()
					kv.startReconfiguration(newConfig)
				} else {
					kv.mu.Unlock()
				}
			}
		}
		time.Sleep(100 * time.Millisecond)
	}
}

func (kv *ShardKV) allShardsReady() bool {
	for shardID, gid := range kv.currentConfig.Shards {
		// 如果有任何一个分片正在拉取中,就不能升级 Config
		if gid == kv.gid && kv.data[shardID].State == Pulling {
			return false
		}
	}
	return true
}

可以看到上面这个代码,在我当前的实现中,进行了一个判断kv.allShardsReady(),这是因为比如分片S属于A,如果我们在 Config N+1 且需要分片S,但S还没从 A 拉过来,就直接跳到了 Config N+2,我们可能会忘记分片S目前还在 A 手里。组 C 会来找我们要数据,但我们自己也是空的,导致数据链条断裂,甚至丢失数据。 只有确保 Config N 的数据完整了,我们才有资格进入 Config N+1,甚至将来把数据传给别人。但是后续也可以优化,但难度会更高,就先这样做了。

当我们成功拉取到新配置后,一定不要直接应用,在raft中所有新数据都要先经过raft进行共识。当raft共识完成,即从appliar loop中接收到新的config后再应用:

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
case InsertConfigOp:
  kv.notifyShard(index)
  if kv.currentConfig.Num+1 == command.Config.Num {
    oldConfig := kv.currentConfig
    kv.calcShardChange(&kv.currentConfig, &command.Config)

    DPrintf(dConfig, "%v: %d -> %d \n", kv, oldConfig.Num, command.Config.Num)

    kv.currentConfig = command.Config
    kv.currentConfig.Groups = make(map[int][]string)
    for k, v := range command.Config.Groups {
      kv.currentConfig.Groups[k] = v
    }
  }

当配置更新时,需要通过 calcShardChange 函数来对比新旧配置的差异,从而驱动分片的状态机流转。

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
func (kv *ShardKV) calcShardChange(oldConfig *shardctrler.Config, newConfig *shardctrler.Config) {
	oldNum := oldConfig.Num

	kv.historyData[oldNum] = make(map[int]map[string]string)
	// 遍历所有的分片
	for shardID := 0; shardID < shardctrler.NShards; shardID++ {
		// 获取该分片在旧配置和新配置中的归属 GID
		oldGID := oldConfig.Shards[shardID]
		newGID := newConfig.Shards[shardID]

		if newGID == kv.gid && oldGID != kv.gid {
			if _, ok := kv.data[shardID]; !ok {
				kv.data[shardID] = &Shard{}
				kv.data[shardID].KV = make(map[string]string)
			}
			if oldGID != 0 {
				kv.needPulling[shardID] = copyConfig(*oldConfig)
				kv.data[shardID].State = Pulling
			} else {
				kv.data[shardID].State = Serving
			}
		}

		if oldGID == kv.gid && newGID != kv.gid {
			if kv.data[shardID].State == Serving {
				kv.data[shardID].State = BePulling

				kv.historyData[oldNum][shardID] = make(map[string]string)
				for k, v := range kv.data[shardID].KV {
					kv.historyData[oldNum][shardID][k] = v
				}
				DPrintf(dDrop, "%v: config num: %d Shard: %d is dropped\n", kv, oldNum, shardID)
			}
		}
	}

}

该函数主要处理两种核心场景:

  1. 分片转入 :
    • 若发现某分片的新 Owner 是我,且旧 Owner 不是我,说明该分片需要迁移进来。
    • 初始分配: 如果旧 Owner 是 0(集群刚启动),直接置为 Serving
    • 迁移分配: 如果来自其他 Group,置为 Pulling 状态,并记录数据来源(needPulling),等待拉取协程去获取数据。
  2. 分片转出:
    • 若某分片的旧 Owner 是我,但新 Owner 不是我,说明该分片要离开。
    • 状态切换: 将分片标记为 BePulling,表示停止服务但保留数据。
    • 数据快照 (History Snapshot): 这是最关键的一步。我们需要将当前分片的数据进行深拷贝并保存到 historyData[oldNum] 中。这样即使当前节点后续覆盖了该分片的数据,新 Owner 依然可以通过 RPC 拿着旧配置号(oldNum)来拉取这份不可变的“历史快照”。

Shard Move

在实现分片迁移时,一个常见的误区(我第一次写就是这么做的,调bug调了3天,,后来把日志的粒度写的非常细才发现)是在 calcShardChange 检测到配置变更后,立即在同一个逻辑流中发起 RPC 拉取数据。这种“事件驱动”的做法非常脆弱:如果网络抖动导致 RPC 失败怎么办?如果节点在拉取过程中崩溃重启了怎么办?这个“拉取事件”就会彻底丢失,导致分片永远处于不可用状态。

因此,我们需要牢记:分布式系统是基于状态机(State Machine)运作的,而不是基于瞬时事件。

1. 职责分离

我们将逻辑拆分为两部分:

  • 状态标记(State Marking): calcShardChange(在 Apply 协程中)只负责修改分片的元数据状态。将分片标记为 Pulling ,将需要被拉取的分片标记为BePulling
  • 状态执行(State Reconciliation): 启动一个独立的后台协程,不断轮询所有分片的状态。一旦发现某个分片处于 Pulling 状态,它就尝试发起数据拉取任务。

2. 永不停止的“拉取循环”

这个独立协程充当了“状态校准器”的角色。只要分片状态还不是 Serving,它就会不知疲倦地重试,直到成功为止。

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
16
17
18
19
20
21
22
func (kv *ShardKV) shardMonitor() {
	for !kv.killed() {
		if kv.isLeader() {
			kv.mu.Lock()
			DPrintf(dFetch, "%v: detect shard needing pulling, len %d need pulling\n", kv, len(kv.needPulling))
			if len(kv.needPulling) > 0 {
				pullingMap := make(map[int]shardctrler.Config, len(kv.needPulling))
				for shardID, config := range kv.needPulling {
					pullingMap[shardID] = copyConfig(config)
				}
				kv.mu.Unlock()

				for shardID, config := range pullingMap {
					go kv.pullShardTask(shardID, &config)
				}
			} else {
				kv.mu.Unlock()
			}
		}
		time.Sleep(100 * time.Millisecond)
	}
}

当后台监控协程发现某个分片需要拉取时,会启动 pullShardTask。这个函数不仅负责网络通信,更关键的是要保证数据写入的安全性一致性

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
func (kv *ShardKV) pullShardTask(shardID int, config *shardctrler.Config) {
	gid := config.Shards[shardID]

	args := FetchShardArgs{
		ConfigNum: config.Num, // 请求特定版本的数据
		Shard:     shardID,
	}

	servers, ok := config.Groups[gid]
	if !ok {
		panic(fmt.Sprintf("can't find gid{%d} from config.Groups", gid))
	}
	DPrintf(dFetch, "%v: is fetching shard %d\n", kv, shardID)

	for kv.isLeader() {
		reply := kv.getShardDataFromGroup(args, servers)

		if reply == nil || reply.Err != OK {
			time.Sleep(100 * time.Millisecond)
			continue
		}

		kv.mu.Lock()

		if kv.currentConfig.Num > config.Num+1 || kv.data[shardID].State != Pulling {
			kv.mu.Unlock()
			DPrintf(dFetch, "Stop pulling shard %d for config %d, current is %d (outdated) or already served/ready", shardID, config.Num, kv.currentConfig.Num)
			return // 直接放弃任务,退出循环
		}

		kv.mu.Unlock()

		op := InsertShardOp{
			ShardID:     shardID,
			ConfigNum:   config.Num + 1,    //这是为了变成新配置准备的数据
			ShardData:   reply.Data,        // 对方发来的 KV
			LastCommand: reply.LastCommand, // 对方发来的去重表
		}

		index, isLeader := kv.submitToRaft(op)
		if !isLeader {
			return
		}
		DPrintf(dFetch, "%v: is raft shard %d, index is %d\n", kv, shardID, index)
		if kv.waitRaft(index) == OK {
			return
		}

		time.Sleep(100 * time.Millisecond)
	}
}

主要包含三个核心步骤:

  1. RPC 拉取: 根据保存的旧配置(config),向分片的前任持有者(Group)发送 FetchShard RPC。
  2. 二次状态检查 : 这是极易被忽略的关键点。 RPC 调用是一个耗时的网络操作,在请求返回期间,本地节点的配置可能已经更新,或者该分片已经被其他协程处理完了。因此,在拿到数据准备提交 Raft 之前,必须加锁再次检查:当前配置是否依然匹配?分片状态是否依然是 Pulling?如果状态变了,必须直接丢弃这份数据,防止“旧数据覆盖新数据”。
  3. 通过 Raft 达成共识 : 拉取到的不仅仅是 KV 数据,还有去重表(LastCommand / Session)。我们不能直接修改内存,而是将数据封装成 InsertShardOp 操作提交给 Raft。

作为数据的提供方,当收到 FetchShard 请求时,节点需要返回指定配置版本的分片数据。

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
func (kv *ShardKV) FetchShard(args *FetchShardArgs, reply *FetchShardReply) {
	if !kv.isLeader() {
		reply.Err = ErrWrongLeader
		return
	}

	kv.mu.Lock()
	defer kv.mu.Unlock()

	shardData, exists := kv.historyData[args.ConfigNum][args.Shard]
	if !exists || shardData == nil {
		// 如果我这儿没有这个分片的数据 (可能已经被 GC 了,或者从未拥有过)
		reply.Err = ErrNoKey
		return
	}

	// DPrintf(dFetch, "%v: the shard %d fetch ask pass shard exist check\n", kv, args.Shard)

	// 深拷贝 KV 数据 (Deep Copy Data)
	reply.Data = make(map[string]string)
	for k, v := range shardData {
		reply.Data[k] = v
	}

	// 深拷贝去重表 (Deep Copy Duplicate Table)
	reply.LastCommand = make(map[int64]int)
	if lastOps, ok := kv.lastCommand[args.Shard]; ok {
		for clientID, seq := range lastOps {
			reply.LastCommand[clientID] = seq
		}
	}

	reply.Err = OK
}

当接收到raft传来的分片后,我们就可以进行接收了

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
case InsertShardOp:
// 幂等性检查
if command.ConfigNum == kv.currentConfig.Num && kv.data[command.ShardID].State == Pulling {
  shardID := command.ShardID

  DPrintf(dFetch, "%v: shardID %d pulling apply success \n", kv, shardID)
  kv.data[shardID].KV = make(map[string]string, len(command.ShardData))
  for k, v := range command.ShardData {
    kv.data[shardID].KV[k] = v
  }

  if _, ok := kv.lastCommand[shardID]; !ok {
    kv.lastCommand[shardID] = make(map[int64]int)
  }
  for k, v := range command.LastCommand {
    if kv.lastCommand[shardID][k] < v {
      kv.lastCommand[shardID][k] = v
    }
  }
  // 状态切换
  kv.data[shardID].State = Serving
  delete(kv.needPulling, shardID)
}

lab5B

200次顺利通过。