You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

6.824 Raft实现applyCh死锁问题排查求助

6.824 Raft实现卡在TestSnapshotAllCrash2D测试的apply阶段

我的6.824 Raft实现通过了所有前置测试,但卡在TestSnapshotAllCrash2D测试的apply阶段,程序出现卡顿。相关日志如下:

...
[Term 1 Server 2] starting new command in 11 for cmd{2942042800837528893}.
[Term 1 Server 2] replicate to 1 success, match index is 10
[Term 1 Server 2] commit index update to 10
[Term 1 Server 2] replicate to 0 success, match index is 10
[Term 1 Server 1] update commit index to 10 from leader.
[Term 1 Server 2] replicate to 1 success, match index is 11
[Term 1 Server 0] update commit index to 10 from leader.
[Term 1 Server 2] commit index update to 11
[Term 1 Server 2] replicate to 0 success, match index is 11
[Term 1 Server 1] updating last applied from 1 to 10
[Term 1 Server 0] updating last applied from 1 to 10
[Term 1 Server 2] updating last applied from 1 to 11

我推测消息已经发送到applyCh但未被接收,但不清楚具体原因。理论上每次调用Start(cmd)都会执行applyCh的消息接收操作。

相关代码片段

1. applier goroutine实现

// go applier() in Make()
func (rf *Raft) applier() {
    // update last applied
    for rf.killed() == false {
        time.Sleep(cpuGap)
        rf.lock("applier")
        if rf.commitIndex > rf.lastApplied {
            ColorPrintf("[Term %d Server %d] updating last applied from %d to %d", rf.currentTerm, rf.me, rf.lastApplied, rf.commitIndex)
        }
        for rf.commitIndex > rf.lastApplied {
            rf.lastApplied++
            rf.applyCh <- ApplyMsg{
                CommandValid: true,
                Command:      rf.accessLog(rf.lastApplied).Command,
                CommandIndex: rf.accessLog(rf.lastApplied).Index,
            }
        }
        rf.unlock("applier")
    }
}

2. 日志复制与commitIndex更新实现

// append entries rpc sender
func (rf *Raft) replicateOneRound(server int) {
    rf.lock("replicateOneRound")
    if rf.state != LEADER {
        rf.unlock("replicateOneRound")
        return
    }

    prevLogIndex := rf.nextIndex[server] - 1

    // if last log index >= nextIndex for a follower: send rpc with log entries starting at nextIndex
    var entries []Entry
    if rf.lastLog().Index >= rf.nextIndex[server] {
        entries = rf.log[rf.nextIndex[server]-rf.SnapshotIndex:]
    } else {
        entries = make([]Entry, 0)
    }

    request := AppendEntriesArgs{
        Term:         rf.currentTerm,
        LeaderId:     rf.me,
        PrevLogIndex: prevLogIndex,
        PrevLogTerm:  rf.accessLog(prevLogIndex).Term,
        Entries:      entries,
        LeaderCommit: rf.commitIndex,
    }
    rf.unlock("replicateOneRound")

    reply := AppendEntriesReply{}
    if rf.sendAppendEntries(server, &request, &reply) {
        rf.lock("sendAppendEntries")

        if rf.currentTerm < reply.Term {
            ColorPrintf("[Term %d Server %d] replicate to %d fail due to a bigger term %d", rf.currentTerm, rf.me, server, reply.Term)
            rf.setState(FOLLOWER, reply.Term)
        }

        if rf.currentTerm == request.Term && rf.state == LEADER {
            if len(request.Entries) != 0 {
                if reply.Success {
                    // append rpc ok
                    rf.matchIndex[server] = request.Entries[len(request.Entries)-1].Index
                    rf.nextIndex[server] = rf.matchIndex[server] + 1
                    ColorPrintf("[Term %d Server %d] replicate to %d success, match index is %d", rf.currentTerm, rf.me, server, rf.matchIndex[server])

                    // update commitIndex
                    rf.tryUpdateCommitIndex()

                } else {
                    // decrement
                    rf.nextIndex[server] = reply.ConflictIndex
                    ColorPrintf("[Term %d Server %d] replicate to %d fail, decrement next index to %d", rf.currentTerm, rf.me, server, rf.nextIndex[server])
                }
            }
        }
        rf.unlock("sendAppendEntries")
    }

}

func (rf *Raft) tryUpdateCommitIndex() {
    match := make([]int, len(rf.peers))
    copy(match, rf.matchIndex)
    sort.Sort(sort.Reverse(sort.IntSlice(match)))
    N := match[len(rf.peers)/2]
    for N > rf.commitIndex {
        if rf.accessLog(N).Term == rf.currentTerm {
            ColorPrintf("[Term %d Server %d] commit index update to %d", rf.currentTerm, rf.me, N)
            rf.commitIndex = N
            break
        } else {
            N--
        }
    }
}

目前发现修改config.go中的applyCh := make(chan ApplyMsg)为applyCh := make(chan ApplyMsg, 1000)后测试可以通过,但希望在不修改config.go的前提下解决问题。

内容的提问来源于stack exchange,提问作者sakamoto

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.26 04:10:45