Raft分布式一致性协议:从入门到崩溃

在AWS躺平地写了几年业务代码之后,感觉还是有必要加深程序员的自我修养,好好学习一下分布式系统到底是怎么工作的,于是搞定了MIT 6.828之后就继续入坑了大名鼎鼎的6.824(现为6.5840)——用golang徒手撸一个Raft协议,并在此基础上构建一个分布式K/V存储系统。写完之后只想感慨这种lab让大学生来做,难度也太离谱了吧,分布式系统的bug调试起来简直是折寿……总之最后还是做到了基于Raft的三个lab,所有测试用例并行压测10k次全都pass,不过受限于课程要求,代码不能公开,于是写篇博客记录一下心得感想。

这篇文章分为无剧透、微剧透以及大量剧透三个部分;第一部分无剧透,介绍有助于避坑的经验感想与调试技术,可以放心阅读。

通关感想

  • 多年来的工作经验只教会我一件事情:编程的本质困难在于你不知道哪里会踩坑。但是做这种精心设计的lab的最大好处就是作者会(尽量)帮你避坑,所以动手写代码之前先熟读所有建议,确保与自己的想法一致,不要去拍脑袋设计一些很fancy的结构,最后发现给自己埋了大雷。
  • 老程序员依然会忽略一些会导致racing condition的细节……如果代码里有细微的bad taste,相信我,它一定会让你经历痛不欲生的debug。
  • goroutine的生命周期有时候跟你理解的不一样,如果出现了泄露,可能会以几乎无法理解的方式影响系统运行,例如定时器(time.Sleep)无法准确执行,导致频繁election timeout从而无法达成共识。如果对此有所怀疑,可以用这里提供的方式,打印出某个时刻所有正在运行的goroutine并加以检查(注意设置debug=2来获得更加可读的stack trace)。总之,当一个Raft instance被Kill()之后,最好保证它所对应的goroutine能够在有限时间内执行结束,哪怕lab并不要求这一点。
  • TA提供的dtest脚本能够方便我们以高并行度压测系统,但在CPU高负载下golang runtime同样会出现定时器不准确的问题,导致Raft运行效率低,从而fail一些test case,这是正常现象,不必过于强求。例如lab4的TestJoinLeave这个case会直接sleep一秒钟来等待某个replica group把shards数据发送出去,随后便会断开它的网络连接,并没有考虑系统高负载的情况下来不及完成发送,最终造成死锁超时。
  • 用channel的时候要特别小心:这个类型是无法被RPC序列化的,因此对它进行反序列化会得到nil,而向nil写入数据则会直接导致阻塞!这个真的是巨坑,难道不应该报runtime error吗……

调试日志与可视化

6.824这门课的lab本质上就是强迫你学习分布式系统的调试技术,因此如何高效打log是其中非常重要的一环,TA甚至专门写了一篇博客来介绍相关技巧,注意其中提供了一个名为dtest的脚本非常有用!基于过去的工作经验,我选择了使用Chromium项目的Perfetto性能数据可视化工具来辅助调试Raft,最终实践下来感觉场景还是比较合适的。这个可视化工具目前仍然对Chromium早年使用的Trace Event Format数据格式提供支持,你只需要按照要求的字段将系统中的关键数据写到json里,就可以直接从web UI读取文件并进行可视化,并且还支持时间线缩放、标记等辅助功能:

图中的每个thread可以对应于Raft server node,每个process则可以对应到lab4里的replica group,所以这个现有的层级结构就非常适合展现Raft的执行细节。我其实只使用了两种event:Duration Events用来可视化每个server在时间线上所处的角色状态;以及Instant Events用来对发生的事件进行标注,把(事件发生后)很详细的内部状态信息以json格式记录下来,非常有助于调试。

可以简单通过环境变量传入要保存的json文件名,并且修改TA提供的dtest脚本,在test failed的时候也同时保存相应的json文件用于后续调试分析。如果记录的事件类型太多,也可以快速撸一个简单的python脚本来对json内容进行过滤,或者是在代码里提供开关等等,方法可以灵活多样。记录事件的代码示意如下(这里比较糙快猛使用了全局变量):

func TraceInstant(name string, server int, group int, timestamp int64, args map[string]any) {
    if !FlagTrace {
        return
    }
    gMutex.Lock()
    defer gMutex.Unlock()
    if gFile == nil {
        return
    }
    logitem := map[string]any{
        "name": name,
        "ph":   "i",
        "pid":  group,
        "tid":  server,
        "ts":   timestamp - gStart,
        "args": args,
    }
    data, _ := json.Marshal(logitem)
    gFile.Write(data)
    gFile.WriteString(",\n")
}

由于要大量访问Raft的内部状态信息,当然也要封装一个方法:

func merge(maps ...map[string]any) map[string]any {
    merged := make(map[string]any)
    for _, m := range maps {
        for k, v := range m {
            merged[k] = v
        }
    }
    return merged
}
func (rf *Raft) GetTraceState() map[string]any {
    return merge(rf.log.GetTraceState(), map[string]any{
        "GID":              rf.getGID(),
        "raft.currentTerm": rf.currentTerm,
        "raft.commitIndex": rf.commitIndex,
        "raft.lastApplied": rf.lastApplied,
        "raft.nextIndex":   fmt.Sprintf("%v", rf.nextIndex),
    })
}
// when recording events:
TraceInstant("NewElection", rf.me, rf.getGID(), time.Now().UnixMicro(), rf.GetTraceState())

很多时候测试会挂在我们代码里的panic上面,最好也把它封装一下,这样就可以在可视化日志里找到相应的事件,并检查它前后都发生了什么:

func (rf *Raft) TracePanic(msg string, context map[string]any) {
    TraceInstant("Panic", rf.me, rf.getGID(), time.Now().UnixMicro(), merge(rf.GetTraceState(), context))
    panic(msg)
}

另一方面,我们需要根据论文的figure 4状态转移图,正确地创建Duration Events来记录每个server在时间线上处于什么role,这样就使得整个可视化更加清晰。示意代码如下:

func (rf *Raft) SwitchToCandidate() {
    now := time.Now().UnixMicro()
    if rf.role == Leader {
        panic("Leader can not become candidate")
    }
    if rf.role == Follower {
        TraceEventEnd(Follower.String(), rf.me, rf.getGID(), now, nil)
    }
    if rf.role != Candidate {
        TraceEventBegin(Candidate.String(), rf.me, rf.getGID(), now, rf.GetTraceState())
    }
    rf.role = Candidate
    // ...
}
func (rf *Raft) SwitchToFollower() { /* ... */ }
func (rf *Raft) SwitchToLeader() { /* ... */ }

注意:在初始化json文件的时候,需要记录一个ts=0的初始事件(SystemStart),这样才能让web UI上显示的时间戳与文件中记录的时间戳一致,不然它会把第一个发生的事件移动到时间零的位置,后续的所有时间戳都会加上这个偏移,导致不方便根据时间戳在UI上查找事件。

以下是我在完成所有labs之后,代码里打印的有助于debug的事件,文件体积允许的情况下最好把状态context记录得详细越好:

  • SystemStart
  • StartCommand
  • NewElection
  • Vote/GotVote
  • Commit/Apply
  • Snapshot/StateMachineSnapshot
  • Heartbeat/AppendEntries
  • InstallSnapshot
  • SendShard/SendShardFailed

实际使用下来,感觉这套调试方案还是有以下缺点:

  • 无法快速搜索一个特定时间点上,或者是满足特定条件的事件,需要手动缩放寻找
  • 事件之间没有办法做关联与快速跳转

⚠️ 以下进入微剧透内容,主要介绍Raft论文以及lab guidance中语焉不详部分的实现思路,以及一些需要注意的点。

思路与细节

Lab2: Raft协议

Election Timeout

If a follower receives no communication over a period of time called the election timeout, then it assumes there is no viable leader and begins an election to choose a new leader.

这里需要注意的是,follower收到任何一种RPC都会重置election timeout定时器,这是因为如果不这样的话,可能会出现刚刚选出leader,又有某个follower timeout之后发起新一轮选举,从而影响稳定性。

另外要注意初始代码里提供的election timeout是偏低的,在TestFigure8Unreliable2C这个case里,由于网络随机抖动延迟设置得较高,并发压测时会出现小概率的超时错误,因为本身网络就不稳定,偏低的timeout会导致当成功选出了leader后,也很容易因为follower没有及时收到心跳而重新发生选举,最后迟迟无法commit log entry。我被这个问题困扰了好一阵子,分析了很久的log,最后确定这不是我的代码实现问题,参考了别人博客之后发现只需要简单将election timeout设置为400-800ms就搞定了……

PS:这里踩的坑就是因为理论基础不过关而导致的,如果已经读过了DDIA,就会意识到Raft的liveness属性是要靠消息处理的有界延迟来保障的,如果这个延迟选得不合适,就会导致系统状态无法推进。

voteFor如何更新?

根据论文figure 4状态转移我们知道,当节点为leader时,该字段没有意义;当节点为candidate时,voteFor显然设置为自己,那么仅当节点状态变为为follower时,才需要设置为null,因为此时还不知道向谁投票。初始化时所有节点均为follower,因此所有voteFor也都是null。

Commit策略

Only log entries from the leader’s current term are committed by counting replicas.

直观上考虑,我们希望一个entry被分发到多数节点之后就可以commit,但如果leader在commit之前就挂掉了呢?新的leader是否可以直接代替上一个leader来完成某个entry的commit呢?论文告诉我们这样的策略是不完善的,举出的反例就是Figure 8——哪怕某个log entry已经被分发到了多数节点,新的leader还是有可能会覆盖掉它们。那么正确的做法是加一个条件,如果想要commit,就要一鼓作气commit到当前term才行,这样也就自然包括了先前term里未commit的entries;但如果做不到,就不能只commit先前term已经分发的replicated log entries。

由于这个策略,当我们debug的时候也需要注意,在term change之后,如果没有新的command进来,就会导致即便已经分发完成的commands也无法commit,这是正常现象。

Snapshot实现细节

将log数组封装为数据结构的时候,极易漏掉一些细节的边界条件,因此要在各种有可能访问越界的地方都加上panic并且打印一些具体的信息,从而有助于修复这类bug。

避免状态机产生倒退

Take care that these snapshots only advance the service's state, and don't cause it to move backwards.

lab里的这句话指出了一个可能导致lab3出bug的细节:当接收InstallSnapshot RPC的时候,如果收到的snapshot包含的日志少于自身当前的日志长度,但是又大于自己的上一个snapshot,我们是会直接保存这个snapshot,但却不一定要apply it to service (state machine)——此时有可能存在rf.commitIndex >= args.LastIncludedIndex,如果apply snapshot就会导致service的状态倒退。我们可以在KV service中维护一个lastAppliedIndex从而避免这种情况,也能通过测试但这不是正确的实现,还是应该从源头保证Raft协议就不会发出错误的apply snapshot。

正确回退nextIndex

这是另外一个容易踩坑的细节,当我们在lab 2C中实现了nextIndex回退优化之后,可能会在AppendEntries RPC里写出这样的逻辑:

if rf.log.IsTermValidAt(args.PrevLogIndex) {
  // ...
} else {
  reply.XLen = rf.log.Length()
  reply.XTerm = -1
    reply.XIndex = -1
    reply.Success = false
  return
}

但当lab 2D中引入日志压缩后,在PrevLogIndex 这个位置找不其Term可能会是两种原因:follower的日志太短,或者这个位置的日志已经被压缩。在后一种情况下,我们需要返回的是第一个尚未进入snapshot的log entry的位置,从而让leader从这个位置开始发送AppendEntries;否则的话,可能会导致leader永远无法把日志成功分发给这个follower。我直到lab4末尾才因为这里的bug触发了错误,实在是隐蔽至极。

按序Apply Message

当我们更新commitIndex之后,如何将ApplyMsg发送到applyCh是一个略有些tricky的地方,尤其是当我们实现了Snapshot调用之后可能会发现拿不到锁而卡死的情况,所以这里最好是异步发送,以免applyCh产生了阻塞,导致Raft上的锁无法释放,卡死整个系统;其次我们也不能每次顺手起一个新的goroutine,这样无法保证写入channel的command index顺序,而是要用一个缓冲队列来处理。InstallSnapshot RPC产生的ApplyMsg也要写入同样的缓冲区,从而在发送给状态机时保证顺序正确。

Lab2是我耗时最久的,而且直到做完lab4,才终于调试到所有case都稳定通过:


⛔️ 以下进入重度剧透,包含了较多需要自己设计推敲的实现细节,阅读下面内容会降低你自己对这些问题的思考深度!

Lab 3: K/V Server

当我们实现了核心Raft协议之后,接下来就要理解它的使用方法,其实本质上就是用来线性化地提交对状态机的并发修改(和读取),从而实现强一致性。这个lab需要考虑清楚的几个重点问题:

  • 如何让client RPC以同步的方式,优雅地从状态机线程拿到执行结果?
  • 理解request deduplication到底是怎么工作的,以及如何优化空间复杂度?
  • 来自client的操作通过rf.Start()成功提交之后就会休眠等待,但并不是每一个成功提交的操作都会最终被commit,根据论文的figure 8我们知道有可能在同一个index最终commit了不同的操作,那么如何优雅地fail掉那些还在等待的请求?

对于第一点,简单分享一下我的实现思路:只需要在Op结构体中保存一个buffered string channel,用于接收从状态机线程执行过程中发送回来的结果即可。

type Op struct {
    Key       string
    Value     string
    Op        string
    ResultCh  chan string
    From      int
    // ...
}
func (kv *KVServer) Get(args *GetArgs, reply *GetReply) {
  // ...
  resultCh := make(chan string, 1)
  index, _, isLeader := kv.rf.Start(Op{
        Op:        GetOp,
        Key:       args.Key,
        ResultCh:  resultCh,
        From:      kv.me,
    // ...
    })
  if isLeader {
    select {
    case result := <-resultCh:
      reply.Value = result
      reply.Err = OK
    }
    // ...
  } else {
    reply.Err = ErrWrongLeader
  }
}

// goroutine StateMachineExecutor
// ...
if op.From == kv.me && op.ResultCh != nil {
  op.ResultCh <- result
}

由于ResultCh在序列化的时候自动被转换成nil,不会被AppendEntries分发到follower上,于是保证只有在接收处理client RPC并阻塞等待的那个server node上,才会需要状态机线程向channel写入数据来解除阻塞。

最后还是免不了一通debug之后总算过关:

Lab 4: Shard K/V Server

这个lab基本上就是只提要求,自己去想办法实现,自由度可以说远高于先前的labs。

几个提示:

  • 在开始写Shard K/V Server之前,一定要用dtest压测Shard Controller模块,保证它是正确的!
  • 先通读lab guide,直接把challenge问题也一起纳入设计范围,没有必要分开去做。

Re-configuration的实现

首先卡住我的一个点在于:当config更新之后,如何将信息同步地传递给所有replica groups?通过lab3我们学习到,分布式系统本质上就是把所有对状态机的并发修改都通过Raft进行线性化提交,所以这里肯定会有一个线程去不断pull shard controller,那么接下来问题是如何把收到的config change应用到Raft。这里有一个关键点思考点:我们能跳过某个或者几个版本的config,直接应用下一个吗?答案是不可以!因为如果允许的话,每个server node都无法保证自己的config更新历史与其他人一致,所以一定会导致某种程度的死锁:例如某个replica group在等待其他group发来自己所需要的shards data,但是相应的group跳过了这个config,所以系统的状态始终无法前进。

因此,在处理re-configuration的时候,同一个replica group内的所有server nodes都要对config变更达成共识,这样才能对外发送状态一致的shards data,这也就意味着我们需要把ConfigChange本身也作为一个log entry扔给底层的raft协议,并且状态机在处理这个command的时候,要求config num必须按序增加,否则就忽略掉。Config puller线程的示意代码如下:

func (kv *ShardKV) configPuller() {
    for !kv.killed() {
        newConfig := kv.mck.Query(-1)
        activeConfig := kv.config.Load()
        if newConfig.Num != activeConfig.Num {
            for num := activeConfig.Num + 1; num <= newConfig.Num && !kv.killed(); num += 1 {
                cfg := kv.mck.Query(num)
                kv.rf.Start(Op{
                    Op:        ConfigChange,
                    NewConfig: &cfg,
                    // ...
                })
            }
        }
        time.Sleep(time.Millisecond * 100)
    }
}

注意在以上代码中我们直接调用了kv.rf.Start(),直观上似乎这样能保证操作一定会进入Raft并由当前leader分发到所有follower nodes,但其实是保证不了的!原因在于有可能每个server node在发起这个操作的时候都不是leader,最终导致这个操作被丢弃掉了;当某个调用返回了isLeader=true的时候,才能保证这个操作进入了Raft,但却仍然不能保证它一定会被执行;只有当操作被发到applyCh才能够真正保证它被执行了。由于上述puller逻辑会读取当前的config num并不断重试,这样才不会影响整体上的正确性,反正就是我只管发,你状态机能不能执行是另说,但总归不会漏。

Request Deduplication

Lab4的dedup逻辑整体与lab3差不多,只是需要注意当我们把shards data发送给其他replica group的时候,也需要把dedup table中相应的元素一起发送过去,这样才能够避免新的replica group从同一个client处接受重复请求。我们可能会对lab3的逻辑简单进行修改,还是用client ID作为key,把shard ID放进value里,于是就可以遍历dedup table,找出并删除那些要发送给其他replica group的shard,这样实现看似没毛病,但却忽略了一个很重要的细节:删掉某个client ID之后,当前replica group如果再收到来自它的请求,还拿什么去做dedup呢?所以这里我们需要的不是简单删除,而是更类似“回退”,实现的办法也很简单——将client ID与shard ID结合起来作为key就好了:

type DedupKey struct {
    ClientId int64
    ShardId    int
}
type DedupEntry struct {
    SeqNumber int32
    Value     string
}
type ShardKV struct {
  //...
  dedup  map[DedupKey]DedupEntry
}
// when sending shards to other replica group:
for key, entry := range kv.dedup {
  for _, shard := range shards {
    if key.Shard == shard {
      // copy the entry into send shard request
      // ...
      delete(kv.dedup, key)
    }
  }
}

这样虽然导致dedup table的体积膨胀,但当我们删除某个(ClientId, ShardId)的时候,仍然保留着这个client ID与余下的shards所组成的keys,从而能够正确对来自这个client的请求进行dedup处理。这里如果没有正确实现的话,会无法通过linearizability测试,极其难以debug。

经过一番痛苦的debug之后,终于稳定通关了lab4所有测试:

Challenge 1 Garbage collection

这里的思考点在于:发送shards data时,必须保证这些数据已经被持久化到相应的replica group里,才能够安全从自身删除。我选择了在处理完ConfigChange操作之后,shards data保存到缓冲区并异步发送,并且将缓冲区也作为persistent state保存起来,但是在这个challenge题目中就会遇到问题——当执行了ConfigChange之后立刻就会触发snapshot,replica group里每个server node缓冲区里的shards data还没有来得及发送就被存了下来,这样就造成了冗余,需要某种机制来确保当shards发送完成之后,要重新做一次snapshot。为了100%通过这个test case,我选择了一个tricky workaround:当每次异步发送shards完成之后,就通过applyCh发送一个特殊的NOP command,将其index设置为-1,来强制触发一次snapshot,这样就可以把上一次写进snapshot的冗余数据给覆盖掉。逻辑大概是这样:

// goroutine SendShards
// ...
if allShardsSent {
  kv.applyCh <- raft.ApplyMsg{CommandValid: true, CommandIndex: -1, Command: Op{
    Op:   Nop,
  }}
}

// goroutine CommandExecutor
// ...
if kv.maxraftstate > 0 && kv.persister.RaftStateSize() >= kv.maxraftstate {
  if cmd.CommandIndex != -1 {
    kv.Snapshot(cmd.CommandIndex)
  } else {
    kv.Snapshot(kv.lastAppliedIndex)
  }
}

虽然这个实现并不好看而且拉低了执行效率,不过这样就可以稳定通过Challenge 1 test case了。

What's Next?

  • 实现线性一致性读:你会发现也能把读路径优化到承载高吞吐,有助于进一步认清同步复制 vs 共识复制之间的 tradeoff
  • 支持分布式事务:获得对 2PC 更加本质的理解,你会意识到一切“原子性”保障的背后都有它的身影

(全文完)🎉

PS:如果你有耐心看到这里,并对我的实现感兴趣,可以email找我索取代码。仅供参考学习,请勿公开发表!

comments powered by Disqus
Published:
2024-02-01
Last modified:
2026-05-20
分类:
Tag: