在作者实现的基于raft算法的分布式数据库中,遇到的一个比较棘手的问题就是在高并发场景下(200 个客户端同时向数据库发送读写请求),日志复制的AppendEntriesRPC中会大量的出现gRPC DeadlineExceeded,即RPC超时的问题,导致follower的日志跟不上leader的日志,leader不断的重试发送RPC。且数据越积累越多,重试的情况也越多,性能也就急剧下降。作者再解决该问题时,也收获了很多的知识,因此写下这篇文章总结分析遇到问题的解决思路。

1 问题出现原因

AppendEntriesRPC调用中错误代码的位置如下。直接原因就是timeout := 200 * time.Millisecond 设置的超时时间200ms,由于高并发下所需要同步的日志过多,grpc超时时间设置不合理。导致频繁报错。

unc (rf *Raft) sendAppendEntries(server int, args *AppendEntriesArgs, reply *AppendEntriesReply) bool {
	if rf.me == server {
		return false
	}

	// 创建新的 gRPC 参数对象,不使用对象池
	grpcEntries := convertToGrpcLogEntries(args.Entries)

	// 注意:LeaderIP 可能为空,需要检查
	var leaderIPBytes []byte
	if args.LeaderIP != "" {
		leaderIPBytes = []byte(args.LeaderIP)
	}

	grpcArgs := &G_AppendEntriesArgs{
		Term:         int64(args.Term),
		LeaderId:     int64(args.LeaderId),
		PrevLogIndex: int64(args.PrevLogIndex),
		PrevLogTerm:  int64(args.PrevLogTerm),
		Entries:      grpcEntries,
		LeaderCommit: int64(args.LeaderCommit),
		LeaderIP:     leaderIPBytes,
	}

	timeout := 200 * time.Millisecond // 默认100ms

	// 设置超时上下文
	ctx, cancel := context.WithTimeout(context.Background(), timeout)
	defer cancel()

	grpcReply, err := rf.peers[server].Grpc_AppendEntries(ctx, grpcArgs)

	// 注意:这里我们不再尝试回收对象,让GC自动处理

	if err != nil {
		// 仅在真实日志复制失败时打印错误,避免心跳噪音
		if len(args.Entries) > 0 {
			fmt.Printf("Grpc_AppendEntries发送给 S%d 失败, err=%v, entries=%d\n", server, err, len(args.Entries))
		}
		return false
	}
	reply.Term = int(grpcReply.Term)
	reply.Success = grpcReply.Success
	reply.ConfilictIndex = int(grpcReply.ConflictIndex)
	reply.ConfilictTerm = int(grpcReply.ConflictTerm)

	return true
}
func (rf *Raft) AppendEntries(args *AppendEntriesArgs, reply *AppendEntriesReply) {
	rf.mu.Lock()
	defer rf.mu.Unlock()
	LOG(rf.me, rf.currentTerm, DDebug, "<- S%d, Receive log, Prev=[%d]T%d, Len()=%d", args.LeaderId, args.PrevLogIndex, args.PrevLogTerm, len(args.Entries))
	reply.Term = rf.currentTerm
	reply.Success = false

	if args.Term < rf.currentTerm {
		LOG(rf.me, rf.currentTerm, DLog2, "<- S%d, Reject log", args.LeaderId)
		return
	}

	if args.Term >= rf.currentTerm {
		rf.becomeFollowerLocked(args.Term)
	}

	// 重置选举超时时间,表示我不再争取作为leader,因为你已经是我的leader
	// 将这个重置放在前面,对 rf.log 进行任何操作之前调用。这确保了在处理完所有事情(例如日志匹配和追加日志条目)后才重置选举计时器
	defer rf.resetElectionTimerLocked()

	// 本地的日志太久没有与leader同步了
	// 检查args.PrevLogIndex是否超出了跟随者的日志长度。如果超出,说明跟随者的日志不够长,拒绝日志追加并记录日志。
	// 需要记录具体的任期和日志中信息进行回复
	if args.PrevLogIndex >= rf.log.size() {
		reply.ConfilictIndex = rf.log.size()
		reply.ConfilictTerm = InvalidIndex
		LOG(rf.me, rf.currentTerm, DLog2, "<- S%d,Reject log,Follower too short,len : %d <= Pre:%d", args.LeaderId, rf.log.size(), args.PrevLogIndex)
		return
	}
	// 检查领导者的前一个日志索引 (PrevLogIndex) 是否小于跟随者的快照最后索引 (snapLastIdx)。
	// 这表示跟随者的日志已经被截断,而该索引对应的条目不再存在于跟随者的日志中。
	if args.PrevLogIndex < rf.log.snapLastIdx {
		reply.ConfilictTerm = rf.log.snapLastTerm
		reply.ConfilictIndex = rf.log.snapLastIdx
		LOG(rf.me, rf.currentTerm, DLog2, "<- S%d, Reject log, Follower log truncated in %d", args.LeaderId, rf.log.snapLastIdx)
		return
	}
	// 代码执行到这里说明此时leader和follower的日志号已经一致
	// 检查任期号是否一致,只有日志号和任期号一致才可以返回成功
	if rf.log.at(args.PrevLogIndex).Term != args.PrevLogTerm {
		reply.ConfilictTerm = rf.log.at(args.PrevLogIndex).Term
		reply.ConfilictIndex = rf.log.firstFor(reply.ConfilictTerm) // 在follower找到这个任期的第一个日志,Leader可以据此快速移动其匹配位置,缩短同步时间。
		LOG(rf.me, rf.currentTerm, DLog2, "<- S%d,Reject log,Pre Log not match,[%d]: T%d != T%d", args.LeaderId, args.PrevLogTerm, rf.log.at(args.PrevLogIndex).Term)
		return
	}
	// 记录leaderIP
	rf.LeaderIP = args.LeaderIP

	// 追加日志
	rf.log.appendFrom(args.PrevLogIndex, args.Entries)
	rf.persistLocked()
	reply.Success = true
	LOG(rf.me, rf.currentTerm, DLog2, "Follower append logs: (%d, %d]", args.PrevLogIndex, args.PrevLogIndex+len(args.Entries))

	// TODO:LeaderCommit
	// 如果leader已提交的日志号大于本地要提交的日志号
	if args.LeaderCommit > rf.commitIndex {
		LOG(rf.me, rf.currentTerm, DApply, "Follower update the commit index %d->%d", rf.commitIndex, args.LeaderCommit)
		rf.commitIndex = args.LeaderCommit   // 将本地的要提交日志号更新
		if rf.commitIndex >= rf.log.size() { // 如果本地提交日志号比本身日志条目还大
			rf.commitIndex = rf.log.size() - 1 // 则将要提交的日志号更新为日志长度 - 1
		}
		rf.applyCond.Signal() // 发送信号准备提交
	}
}

由于gRPC使用的是一元RPC模型,每次leader发送rpc给follower时都必须等待follower的回复。如果follower没有及时处理好RPC,那么就会出现大量超时的问题。那么最重要的就是找到为什么follower没有及时回复leader呢?瓶颈究竟在哪?从代码上看,从目前作者首先想到的是两个地方:

  • 从leader端我们可以通过延迟超时时长,将200ms延长,那么follower就有足够的时间
  • 从follower端看,即看AppendEntries的实现,我们发现主要的瓶颈应该是persistLocked这个磁盘IO的持久化中

伏笔:而且目前整个raft使用的是一把大锁,如果持久化时间过长导致持锁时间长。那么心跳逻辑也会受到影响,也会出现leader无故易主的情况

高并发写请求 → 大量日志生成 → 需要持久化到磁盘
                                    ↓
                         磁盘IO缓慢(瓶颈)
                                    ↓
              ┌────────────────────┴────────────────────┐
              ↓                                         ↓
         持久化阻塞                              资源被占用
              ↓                                         ↓
         锁持有时间过长                          心跳无法及时发送
              ↓                                         ↓
         其他操作等待锁                          Follower超时触发选举
              ↓                                         ↓
         RPC超时错误                              Leader频繁更换

优化1 批量异步持久化

既然是持久化时间过长导致的问题,因此作者想到的第一个解决方案就是进行批量异步的持久化。即将持久化与RPC逻辑分开处理

目前请求的链路

正常路径(低并发):
  Leader 发送 AppendEntries → Follower 收到 → 处理 → 返回 → 总耗时 < 200ms ✅

高并发时的阻塞链:
  客户端 200 并发写请求
    ↓
  batch_manager.Submit() × 200 → 全部阻塞等待 RespCh
    ↓
  rf.Start(ops) → 写入 Raft 日志
    ↓
  startReplication() → 向 S1, S2 发送 gRPC AppendEntries
    ↓
  ⚠️ Follower 端同时收到大量 AppendEntries
    ↓
  Follower: AppendEntries() 加 rf.mu.Lock() 排队等待
    ↓
  ⚠️ rf.mu.Lock() 等待时间 > 200ms → DeadlineExceeded ❌

即每个请求都要一次同步写磁盘,占用时间比较长,处理队列堆积导致RPC过载

改为异步持久化的数据流如下:

Follower 收到 AppendEntries:
  ├── 日志追加到内存 (快)
  ├── persistLockedAsync() → 放入队列,立即返回
  └── 返回 Success 给 Leader

后台 AsyncPersister:
  ├── 每 10ms 或累积 10 条日志
  └── 批量写入磁盘 (减少 IO 次数)

优化2 同步批量持久

但是其实上面异步持久化虽然能够提高性能,且也没有出现了grpc超时的错误。但是却忽略了一个重要的问题——raft安全性的问题。在raft论文规定中,Follower收到日志后必须先持久化,再回复Leader。如果异步处理:

Follower收到日志 → 追加到内存 → 立即回复Leader → 异步刷盘
                                     ↓
                            Leader收到多数确认 → 提交日志
                                     ↓
                            Follower崩溃 → 日志未落盘
                                     ↓
                            新Leader选举 → 数据丢失!
  1. 双Leader场景 :旧Leader提交日志后立即崩溃,新Leader当选
  2. 日志不一致 :崩溃的Follower重启后缺少已提交的日志
  3. 违反Raft安全性 :已提交的日志理论上不应该丢失

因此,将异步的批量改为同步批量。使用channel等待异步处理完成后,使用管道的信号通知处理完成。而后才回复leader成功的信号。那么此时大致同步批量持久化的步骤如下:

┌─────────────────────────────────────────────────────────────────────────┐
│                    批量同步持久化流程                                    │
├─────────────────────────────────────────────────────────────────────────┤
│                                                                         │
│  Follower 收到 AppendEntries                                            │
│       │                                                                 │
│       ▼                                                                 │
│  追加日志到内存                                                          │
│       │                                                                 │
│       ▼                                                                 │
│  提交到 AsyncPersister 队列                                             │
│       │                                                                 │
│       ▼                                                                 │
│  ┌─────────────────────────────────────────────────────────────────┐   │
│  │  batchWorker 积累请求                                            │   │
│  │                                                                   │   │
│  │  请求 1 ──┐                                                       │   │
│  │  请求 2 ──┼──► 批量写入磁盘(一次 fsync)                         │   │
│  │  请求 3 ──┘                                                       │   │
│  │                                                                   │   │
│  │  触发条件:                                                       │   │
│  │  - 达到 batchSize(如 100 条)                                   │   │
│  │  - 超时(如 10ms)                                               │   │
│  │  - 优先级请求                                                    │   │
│  └─────────────────────────────────────────────────────────────────┘   │
│       │                                                                 │
│       ▼                                                                 │
│  close(doneCh) 通知所有等待者                                           │
│       │                                                                 │
│       ▼                                                                 │
│  Follower 收到通知,返回 Success                                        │
│                                                                         │
└─────────────────────────────────────────────────────────────────────────┘

这样的使用同步批量超时持久化的方式,就是符合raft的安全性的。只有在日志持久化成功后才能够回复leader响应的信息。

优化3 同步非批量持久化

但是,通过测试我发现,虽然确实也很少出现了grpc超时的问题。但是问题是不是真是因为批量持久化的效果呢?其实我发现并不是这样的,因为我计算了下面的一个结果:

Follower 接收 AppendEntries 的频率:
─────────────────────────────────────────────

unifiedReplicator 每 80ms 发送一次
↓
Follower 每秒接收 ≈ 12.5 个 AppendEntries
↓
Follower 每秒调用 submitPersistAsync ≈ 12.5 次
↓
每 10ms 内平均到达 ≈ 0.125 个请求

我的日志复制是每隔80ms发送一次的,那么每个follower每秒接收到的是12.5个RPC,那么follower每秒调用批量持久化函数应该是12.5次。但是批量持久化超时时长为10ms,也就是说每10ms就会强制调用一次持久化。也就是每10ms平均到达0.125个请求。以一个时间轴为例子如下:

时间轴(10ms 一个窗口):
─────────────────────────────────────────────

0ms    10ms    20ms    30ms    40ms    50ms    60ms    70ms    80ms
│       │       │       │       │       │       │       │       │
[1]     [0]     [0]     [0]     [0]     [0]     [0]     [0]     [1]
 ↑                               ↑                               ↑
 超时触发                        超时触发                        超时触发
 刷盘(1条)                       刷盘(0条)                       刷盘(1条)

结果:每次刷盘时,批量大小几乎都是 1,极少达到 2!

也就是说当前的批量设置根本没有任何意义,因为根本到达不了批量效果。

那应该如何设置呢?

方案A:增大 batchTimeout(利用批量)

  • batchSize = 5, batchTimeout = 500ms
  • 效果
    • 每 500ms 批量写入 5-6 条日志
    • fsync 次数减少 5-6 倍
    •  但每个请求延迟增加 500ms ❌ 不可接受

方案B:减小 batchTimeout(减少等待)

  • batchSize = 1, batchTimeout = 1ms
  • 效果
    • 相当于每次立即刷盘
    • 几乎没有额外等待延迟
    • 但没有批量优势

方案C:去掉批量,直接同步持久化

  •  不同AsyncPersister,直接调用 persister.Save()
  • 效果
    • 没有额外等待延迟
    • 代码更简单
    • 但每次都要 fsync

当前场景下,批量没有起到作用,反而增加了延时batchTimeOut的10ms固定时长。因此作者决定将批量持久化功能去掉,改为原先的同步持久化。

优化4 减小锁的粒度

从上面讲了那么多的内容,最后其实都没有使用上。重新改回同步非批量持久化又会出现超时的问题。这时其实在前面的伏笔处,我有提到,整个raft目前使用的是一把大锁。我们把这个问题再回到最初的考虑的思路中。作者一味的想着通过减少持久化的时间,从而减少锁的持有时间,从理论上将应该会有效果,但是实际测试似乎没有这么理想。因此,我们可以直接减少锁的持有实际,从而减少排队等待锁的时间,从而减少rpc超时的时间。

func (rf *Raft) AppendEntries(args *AppendEntriesArgs, reply *AppendEntriesReply) {
	rf.mu.Lock()
	LOG(rf.me, rf.currentTerm, DDebug, "<- S%d, Receive log, Prev=[%d]T%d, Len()=%d", args.LeaderId, args.PrevLogIndex, args.PrevLogTerm, len(args.Entries))
	reply.Term = rf.currentTerm
	reply.Success = false

    // ....一些日志验证代码

	rf.log.appendFrom(prevLogIndex, entries)

    // 进行制作持久化的二进制数据但是不落盘
	raftstate, snapshot := rf.preparePersistData()

    // ...

	LOG(rf.me, rf.currentTerm, DLog2, "Follower append logs: (%d, %d]", prevLogIndex, prevLogIndex+len(entries))

    // 提前释放锁
	rf.mu.Unlock()

	// 锁外落盘,同步持久化
	rf.persistDirect(raftstate, snapshot)

	reply.Success = true
	LOG(rf.me, rf.currentTerm, DPersist, "Follower persist completed for logs: (%d, %d]", prevLogIndex, prevLogIndex+len(entries))
	if needSignal {
		rf.mu.Lock()
		rf.applyCond.Signal()
		rf.mu.Unlock()
	}
}

通过测试确实发现,grpc超时的报错减少了,因此真相确实是因为锁的持有时间导致的:

  • 原始代码
    • 持锁做文件 IO(20ms)
    • 其他请求被阻塞
    • 心跳被阻塞 → 可能
    • 超时频繁
  • 修改后:
    • 锁外等待持久化
    • 其他请求不受影响
    • 心跳正常
    • 超时减少

作者此时又想到一个问题,批量持久化理论上应该也能够起到作用。如果我减少心跳的时间,如果将leader发送日志的频率从每80ms一次减少到1ms或者更少,那么批量持久化就可以被启用,是否意味着能够又性能的提升呢?

对比分析:

指标 80ms 间隔 1ms 间隔
RPC 发送频率 12.5/秒 1000/秒
网络传输量 高(80倍)
CPU 消耗 高(序列化/反序列化)
批量持久化效果 无效 有效(10个/批)
fsync 次数 12.5/秒 100/秒(减少10倍)
单请求延迟 10-50ms 10-50ms(不变)
吞吐量 中等 可能更高

理论上可行,但实际上需要考虑RPC和CPU的处理能力

假设Follower处理能力:

  • rpc序列化/反序列化:1-5ms/请求
  • 持久化:10ms-50ms/请求
  • 总处理时间:11ms-55ms/请求

最大吞吐量:

  • 理论上:1000ms/11ms ≈ 90请求/秒
  • 实际(考虑并发):可能100-500请求/秒
  • 1ms间隔要求:1000请求/秒 ❌ 超过处理能力

结果:导致请求会堆积,延迟增加,最终超时

优化5 流水线复制

也就是说一味的提高日志复制的频率从而提高系统的性能,理论上是可行的。但是性能瓶颈会出现在RPC处理能力这边,需要提高RPC的处理能力。ETCD和TIKV有给我们答案——采用流水线复制。

5.1 什么是流水线复制

流水线复制(Pipeline Replication)是一种优化 Raft 日志复制性能的技术。它允许 Leader 在不等待前一个 AppendEntries RPC 响应的情况下,立即发送下一个请求,从而隐藏网络延迟和 Follower 处理时间

传统的串行复制问题:

Leader                    Follower
  │                          │
  │── AE1 ──────────────────→│
  │                          │ 处理 AE1(持久化 50ms)
  │←── Response1 ─────────── │
  │                          │
  │── AE2 ──────────────────→│  ← 必须等待响应才能发送下一个
  │                          │ 处理 AE2(持久化 50ms)
  │←── Response2 ─────────── │
  │                          │
  │── AE3 ──────────────────→│
  │                          │ 处理 AE3(持久化 50ms)
  │←── Response3 ─────────── │

问题:

  • 发送频率收follower处理时间限制
  • 如果 Follower 处理时间 > 80ms,发送频率 < 12.5/秒,会导致
    • 日志积压 ——> 最终超时
Leader                      Follower
  │                            │
  │──── AppendEntries 1 ─────→ │
  │──── AppendEntries 2 ─────→ │  ← 不等响应就发送下一个
  │──── AppendEntries 3 ─────→ │
  │                            │
  │←─── Response 1 ─────────── │
  │←─── Response 2 ─────────── │
  │←─── Response 3 ─────────── │

优势:

  • Leader 发送不受 Follower 处理时间影响
  • 日志积压少、减少超时

5.2 流水线复制是否是安全的

与异步复制区别:

模式 处理流程 安全性 备注
异步持久化(不安全) Follower 接收日志 → 写入内存 → 立即返回 Success 违反 Raft 安全性(崩溃会丢失日志) 不等待持久化即响应
流水线复制(安全) Follower 接收日志 → 持久化 → 返回 Success 符合 Raft 安全性(持久化成功才返回) Leader 不等待响应就发送下一个请求
对比点 流水线复制 异步持久化
Leader 行为 不等待响应就发送下一个请求 (同左,但本质区别在 Follower)
Follower 响应条件 必须持久化成功才返回 不等待持久化即返回
Raft 安全性 安全 不安全

5.3 核心数据结构

1. PipelineReplicator

流水线复制管理器,负责管理所有 peer 的流水线。

type PipelineReplicator struct {
    rf        *Raft            // Raft 实例
    term      int              // 当前任期
    pipelines []*peerPipeline  // 每个 peer 的流水线
    stopCh    chan struct{}    // 停止信号
    running   int32            // 运行状态
}

2. peerPipeline

单个 peer 的流水线,管理 inflight 请求队列。

type peerPipeline struct {
    peer     int                // peer ID
    rf       *Raft              // Raft 实例
    inflight []*pipelineEntry   // inflight 请求队列
    stopCh   chan struct{}      // 停止信号
    cond     *sync.Cond         // 条件变量
    running  int32              // 运行状态
}

3. pipelineEntry

流水线条目

type pipelineEntry struct {
    index     int                // 日志索引
    term      int                // 任期
    args      *AppendEntriesArgs // RPC 参数
    timestamp time.Time          // 创建时间
}

4. inflight 请求

已发送但未收到响应的请求,最多 16 个。

const (
    maxInflight = 16  // 每个 peer 最多 16 个 inflight 请求
)

5.4 整体架构

6 总结

作者感觉后续如果想要继续提高日志复制的性能(防止RPC超时),其实可以从以下方面进行考虑:

  • 锁竞争 —— 高并发gRPC超时
  • 磁盘IO瓶颈 —— Follower持久化慢
  • 内存占用 —— inflight队列过多

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐