Skip to content
Go back

RocketMQ DLedger——基于Raft的CommitLog多数派提交

RocketMQ DLedger:基于 Raft 的 CommitLog 多数派提交

一句话结论(30s)

DLedger 把 Raft 共识协议嵌入 CommitLog 写入路径,用「多数派提交」替代传统异步主从复制,把主库故障切换从手动操作变成自动选举。因为只有多数派节点落盘才 commit 并返回 ACK,所以任意少数节点宕机系统照常工作、已 ACK 消息不丢;代价是每条消息写延迟增加、吞吐约降 10-15%。

核心原理(2min)

Producer 写 Leader → Leader 写本地 DLedger Log(即 CommitLog Entry)后并行向 Follower 发 AppendEntries RPC → Follower 落盘返回 ACK → Leader 收到多数派(含自身)ACK 后 commit 并回 ACK。关键机制:只同步 CommitLog,ConsumeQueue/IndexFile 作为派生索引由各节点本地异步重建(纳入 Raft 会让共识数据量膨胀 3-5 倍);选主时 lastLogIndex 最大的节点当选,Raft 保证新 Leader 拥有所有已 committed 日志;3 节点容 1 故障、5 节点容 2 故障。

底层深入(5-10min)

源码定位:Raft 核心并不在 RocketMQ 主仓,而在独立依赖 io.openmessaging.storage:dledger(0.3.2,包名 io.openmessaging.storage.dledger)。RocketMQ 主仓里属于 DLedger 的只有一层集成代码 DLedgerCommitLogorg.apache.rocketmq.store.dledger)。下面分别看两边的真实源码。

传统主从的问题

非 DLedger 模式的主从复制是异步的——主库写 CommitLog 后立即返回 ACK 给 Producer,再异步同步给从库。主库宕机后,未同步的已 ACK 消息永久丢失,且需人工提升从库为主库。

异步复制的问题到底出在哪? 出在「已 ACK」和「已落盘到从库」之间有时间差。主库写本地 CommitLog 后立刻回 ACK,Producer 以为成功,但消息可能还躺在主库内存或正在异步同步的路上,主库一宕,这批「已 ACK 但未同步」的消息就永久丢失,而且从库能不能顶上、数据新不新全靠人工判断。所以问题不是「复制太慢」,而是「成功语义和真实持久化脱节」——这就是为什么 DLedger 要改用 Raft。

DLedger 的核心改动

DLedger 将 Raft 共识协议嵌入 CommitLog 的写入路径:

Producer → Leader Broker:
  1. Leader 写入 DLedger Log(即 CommitLog Entry)到本地磁盘
  2. Leader 并行向 Follower1, Follower2 发送 AppendEntries RPC
  3. Follower 写本地 DLedger Log → 返回 ACK
  4. Leader 收到多数派(含自身)ACK → commit(更新 commitIndex)
  5. Leader → Producer: 返回写入成功

步骤 2-4 是关键:只有多数派(3 节点中 ≥2 个)写入成功,Leader 才确认提交。任意少数节点宕机,系统继续工作。

RocketMQ 侧,DLedgerCommitLog 构造时把 broker 配置装配成 DLedger 的配置并启动 Raft 内核:

dLedgerConfig = new DLedgerConfig();
dLedgerConfig.setEnableDiskForceClean(defaultMessageStore.getMessageStoreConfig().isCleanFileForciblyEnable());
dLedgerConfig.setStoreType(DLedgerConfig.FILE);
dLedgerConfig.setSelfId(defaultMessageStore.getMessageStoreConfig().getdLegerSelfId());
dLedgerConfig.setGroup(defaultMessageStore.getMessageStoreConfig().getdLegerGroup());
dLedgerConfig.setPeers(defaultMessageStore.getMessageStoreConfig().getdLegerPeers());
dLedgerConfig.setStoreBaseDir(defaultMessageStore.getMessageStoreConfig().getStorePathRootDir());
dLedgerConfig.setDataStorePath(defaultMessageStore.getMessageStoreConfig().getStorePathDLedgerCommitLog());
...
id = Integer.parseInt(dLedgerConfig.getSelfId().substring(1)) + 1;
dLedgerServer = new DLedgerServer(dLedgerConfig);
dLedgerFileStore = (DLedgerMmapFileStore) dLedgerServer.getdLedgerStore();

每条消息真正写盘时,asyncPutMessage 不再像普通 CommitLog 那样「本地落盘即返回」,而是把消息序列化成 byte[] 塞进 AppendEntryRequest,交给 DLedgerServer.handleAppend 走 Raft 共识:

AppendEntryRequest request = new AppendEntryRequest();
request.setGroup(dLedgerConfig.getGroup());
request.setRemoteId(dLedgerServer.getMemberState().getSelfId());
request.setBody(encodeResult.getData());
dledgerFuture = (AppendFuture<AppendEntryResponse>) dLedgerServer.handleAppend(request);
if (dledgerFuture.getPos() == -1) {
    return CompletableFuture.completedFuture(new PutMessageResult(PutMessageStatus.OS_PAGE_CACHE_BUSY, new AppendMessageResult(AppendMessageStatus.UNKNOWN_ERROR)));
}
...
return dledgerFuture.thenApply(appendEntryResponse -> {
    PutMessageStatus putMessageStatus = PutMessageStatus.UNKNOWN_ERROR;
    switch (DLedgerResponseCode.valueOf(appendEntryResponse.getCode())) {
        case SUCCESS:
            putMessageStatus = PutMessageStatus.PUT_OK;
            break;
        case INCONSISTENT_LEADER:
        case NOT_LEADER:
        case LEADER_NOT_READY:
        case DISK_FULL:
            putMessageStatus = PutMessageStatus.SERVICE_NOT_AVAILABLE;
            break;
        case WAIT_QUORUM_ACK_TIMEOUT:
            //Do not return flush_slave_timeout to the client, for the client will ignore it.
            putMessageStatus = PutMessageStatus.IN_SYNC_REPLICAS_NOT_ENOUGH;
            break;
        case LEADER_PENDING_FULL:
            putMessageStatus = PutMessageStatus.OS_PAGE_CACHE_BUSY;
            break;
    }
    ...
});

DLedgerCommitLog 只负责把 dledger 内部的响应码映射回 RocketMQ 语义——尤其 WAIT_QUORUM_ACK_TIMEOUT(等多数派 ACK 超时)被转成 IN_SYNC_REPLICAS_NOT_ENOUGH,而不是伪装成 flush_slave_timeout 让客户端忽略。这一步正是「用 Raft 多数派替代主从异步复制」的落点:异步复制里主库本地落盘就回 ACK,而这里必须等 handleAppend 返回 SUCCESS 才算提交。

💭 思考:为什么 WAIT_QUORUM_ACK_TIMEOUT 要映射成 IN_SYNC_REPLICAS_NOT_ENOUGH,而不是沿用 flush_slave_timeout 让客户端忽略?—— 因为「多数派没凑齐」是真实的可靠性降级,必须让 Producer 感知并重试/告警;如果伪装成一个客户端会忽略的错误码,Producer 会误判成功,反而丢掉「已 ACK 不丢」的保证。所以这一步是把 Raft 的失败语义「如实翻译」给上层,而不是糊过去。

为什么必须等多数派 ACK 才算提交,而不是本地落盘就返回? 因为如果本地落盘就回 ACK,主库宕机后这条消息可能只在旧主库磁盘上,新主(从库)没有,已 ACK 的消息照样丢——又回到异步复制的老问题。只有「多数派落盘 + commit 后才 ACK」,才能保证任意少数节点宕机时,已 ACK 的消息一定存在于某个存活的多数派节点上,选出的新 Leader 一定包含它,从而做到已 ACK 不丢。

日志复制:DLedgerServer.handleAppend

Leader 侧日志复制的入口在 dledger 的 DLedgerServer.handleAppend(0.3.2 版真实源码):

@Override
public CompletableFuture<AppendEntryResponse> handleAppend(AppendEntryRequest request) throws IOException {
    try {
        PreConditions.check(memberState.getSelfId().equals(request.getRemoteId()), DLedgerResponseCode.UNKNOWN_MEMBER, "%s != %s", request.getRemoteId(), memberState.getSelfId());
        PreConditions.check(memberState.getGroup().equals(request.getGroup()), DLedgerResponseCode.UNKNOWN_GROUP, "%s != %s", request.getGroup(), memberState.getGroup());
        PreConditions.check(memberState.isLeader(), DLedgerResponseCode.NOT_LEADER);
        PreConditions.check(memberState.getTransferee() == null, DLedgerResponseCode.LEADER_TRANSFERRING);
        long currTerm = memberState.currTerm();
        if (dLedgerEntryPusher.isPendingFull(currTerm)) {
            AppendEntryResponse appendEntryResponse = new AppendEntryResponse();
            appendEntryResponse.setGroup(memberState.getGroup());
            appendEntryResponse.setCode(DLedgerResponseCode.LEADER_PENDING_FULL.getCode());
            appendEntryResponse.setTerm(currTerm);
            appendEntryResponse.setLeaderId(memberState.getSelfId());
            return AppendFuture.newCompletedFuture(-1, appendEntryResponse);
        } else {
            if (request instanceof BatchAppendEntryRequest) {
                ...
            } else {
                DLedgerEntry dLedgerEntry = new DLedgerEntry();
                dLedgerEntry.setBody(request.getBody());
                DLedgerEntry resEntry = dLedgerStore.appendAsLeader(dLedgerEntry);
                return dLedgerEntryPusher.waitAck(resEntry, false);
            }
        }
    } catch (DLedgerException e) {
        LOGGER.error("[{}][HandleAppend] failed", memberState.getSelfId(), e);
        AppendEntryResponse response = new AppendEntryResponse();
        response.copyBaseInfo(request);
        response.setCode(e.getCode().getCode());
        response.setLeaderId(memberState.getLeaderId());
        return AppendFuture.newCompletedFuture(-1, response);
    }
}

handleAppend 先校验 selfId、group、isLeader 和 term,然后 dLedgerStore.appendAsLeader 只把 entry 追加到本地日志(此时尚未 commit),再交给 dLedgerEntryPusher.waitAck 异步等待多数派 Follower 的 ACK。pending 队列满时直接回 LEADER_PENDING_FULL 拒绝写入,避免内存被未决请求撑爆。关键区别:普通主从模式下主库写本地 CommitLog 后就回 ACK,而这里 ACK 必须等 waitAck 里的多数派(含自身)确认完成——这保证任何少数节点宕机,已返回 SUCCESS 的消息都不会丢。

为什么 ConsumeQueue 不参与 Raft?

ConsumeQueue 和 IndexFile 是 CommitLog 的派生索引——可以从 CommitLog 完全重建。把它们纳入 Raft 同步会让共识数据量膨胀 3-5 倍。

DLedger 的做法:只同步 CommitLog,ConsumeQueue/IndexFile 由每个 Broker 本地异步重建。 Follower 晋升为 Leader 后,需等待这两个索引重建完成才能提供消费服务。

为什么 ConsumeQueue 可以不参与 Raft 同步? 因为它不是源数据,而是 CommitLog 的派生索引——随时可以从 CommitLog 完整重建。把派生数据纳入共识,等于让每次共识都要多同步 3-5 倍的数据量,纯属浪费。所以 DLedger 只同步 CommitLog,索引由各节点本地异步重建,用「可重建性」换「共识开销」。

选主与故障切换

Leader 宕机 → Follower 选举超时 → 发起 RequestVote → lastLogIndex 最大的节点当选(Raft 保证:新 Leader 拥有所有已 committed 的日志)→ 原 Leader 恢复后自动降级为 Follower。

3 节点容忍 1 故障,5 节点容忍 2 故障。切换对 Producer 透明(重试机制自动发现新 Leader)。

为什么选主要比较 lastLogIndex(日志新旧),而不是随便选一个? 因为如果让日志落后的节点当选 Leader,它会把更旧的数据当作权威,导致已经 committed 的消息被覆盖丢失。Raft 用「日志最新者才有资格当选」这条约束,保证新 Leader 一定拥有所有已 committed 的日志,切主不丢已提交数据。

投票的处理逻辑在 dledger 的 DLedgerLeaderElector.handleVote(0.3.2 版真实源码):

public CompletableFuture<VoteResponse> handleVote(VoteRequest request, boolean self) {
    synchronized (memberState) {
        if (!memberState.isPeerMember(request.getLeaderId())) {
            LOGGER.warn("[BUG] [HandleVote] remoteId={} is an unknown member", request.getLeaderId());
            return CompletableFuture.completedFuture(new VoteResponse(request).term(memberState.currTerm()).voteResult(VoteResponse.RESULT.REJECT_UNKNOWN_LEADER));
        }
        if (!self && memberState.getSelfId().equals(request.getLeaderId())) {
            LOGGER.warn("[BUG] [HandleVote] selfId={} but remoteId={}", memberState.getSelfId(), request.getLeaderId());
            return CompletableFuture.completedFuture(new VoteResponse(request).term(memberState.currTerm()).voteResult(VoteResponse.RESULT.REJECT_UNEXPECTED_LEADER));
        }

        if (request.getLedgerEndTerm() < memberState.getLedgerEndTerm()) {
            return CompletableFuture.completedFuture(new VoteResponse(request).term(memberState.currTerm()).voteResult(VoteResponse.RESULT.REJECT_EXPIRED_LEDGER_TERM));
        } else if (request.getLedgerEndTerm() == memberState.getLedgerEndTerm() && request.getLedgerEndIndex() < memberState.getLedgerEndIndex()) {
            return CompletableFuture.completedFuture(new VoteResponse(request).term(memberState.currTerm()).voteResult(VoteResponse.RESULT.REJECT_SMALL_LEDGER_END_INDEX));
        }

        if (request.getTerm() < memberState.currTerm()) {
            return CompletableFuture.completedFuture(new VoteResponse(request).term(memberState.currTerm()).voteResult(VoteResponse.RESULT.REJECT_EXPIRED_VOTE_TERM));
        } else if (request.getTerm() == memberState.currTerm()) {
            if (memberState.currVoteFor() == null) {
                //let it go
            } else if (memberState.currVoteFor().equals(request.getLeaderId())) {
                //repeat just let it go
            } else {
                if (memberState.getLeaderId() != null) {
                    return CompletableFuture.completedFuture(new VoteResponse(request).term(memberState.currTerm()).voteResult(VoteResponse.RESULT.REJECT_ALREADY_HAS_LEADER));
                } else {
                    return CompletableFuture.completedFuture(new VoteResponse(request).term(memberState.currTerm()).voteResult(VoteResponse.RESULT.REJECT_ALREADY_VOTED));
                }
            }
        } else {
            //stepped down by larger term
            changeRoleToCandidate(request.getTerm());
            needIncreaseTermImmediately = true;
            //only can handleVote when the term is consistent
            return CompletableFuture.completedFuture(new VoteResponse(request).term(memberState.currTerm()).voteResult(VoteResponse.RESULT.REJECT_TERM_NOT_READY));
        }

        if (request.getTerm() < memberState.getLedgerEndTerm()) {
            return CompletableFuture.completedFuture(new VoteResponse(request).term(memberState.getLedgerEndTerm()).voteResult(VoteResponse.RESULT.REJECT_TERM_SMALL_THAN_LEDGER));
        }

        if (!self && isTakingLeadership() && request.getLedgerEndTerm() == memberState.getLedgerEndTerm() && memberState.getLedgerEndIndex() >= request.getLedgerEndIndex()) {
            return CompletableFuture.completedFuture(new VoteResponse(request).term(memberState.currTerm()).voteResult(VoteResponse.RESULT.REJECT_TAKING_LEADERSHIP));
        }

        memberState.setCurrVoteFor(request.getLeaderId());
        return CompletableFuture.completedFuture(new VoteResponse(request).term(memberState.currTerm()).voteResult(VoteResponse.RESULT.ACCEPT));
    }
}

handleVote 实现了 Raft 的 RequestVote 语义:先比 ledgerEndTerm、再比 ledgerEndIndex(日志最新者才有资格当选),最后比 term——term 更旧直接 REJECT_EXPIRED_VOTE_TERM,同 term 已投给别人则 REJECT_ALREADY_VOTED,收到更大 term 则主动退为 candidate 并立即触发下一轮。只有日志不落后、term 合法、且本节点在该 term 还没投过票时,才 setCurrVoteFor 并返回 ACCEPT。这套 term + 日志新旧的双重比较,正是 Raft 保证「新 Leader 一定拥有所有已 committed 日志」的机制。

💭 思考:为什么选主要同时过「日志新旧」和「term」两道闸门,而不是只比日志?—— 因为两道闸门各管一件事:日志新旧保证「新 Leader 一定拥有所有已 committed 日志」,数据不丢;term 保证「任期单调递增、不倒退」,防止过期 term 的老 Leader 抢回领导权引发脑裂。只比日志会漏掉「日志一样新、但 term 已过期」的场景,所以 handleVote 里两者缺一不可。

多数派判定与计票逻辑,分别在 MemberState.isQuorumDLedgerLeaderElector.maintainAsCandidate 中:

public boolean isQuorum(int num) {
    return num >= ((peerSize() / 2) + 1);
}
if (alreadyHasLeader.get()
        || memberState.isQuorum(acceptedNum.get())
        || memberState.isQuorum(acceptedNum.get() + notReadyTermNum.get())) {
    voteLatch.countDown();
}
...
} else if (memberState.isQuorum(acceptedNum.get())) {
    parseResult = VoteResponse.ParseResult.PASSED;
} else if (memberState.isQuorum(acceptedNum.get() + notReadyTermNum.get())) {
    parseResult = VoteResponse.ParseResult.REVOTE_IMMEDIATELY;
}
...
if (parseResult == VoteResponse.ParseResult.PASSED) {
    LOGGER.info("[{}] [VOTE_RESULT] has been elected to be the leader in term {}", memberState.getSelfId(), term);
    changeRoleToLeader(term);
}

isQuorum 是多数派判定的唯一实现:num >= peerSize()/2 + 1,即 3 节点需 2 票、5 节点需 3 票。maintainAsCandidate 收集各节点的 VoteResponse,当 acceptedNum 达到多数派时 voteLatch 提前 countDown,最终 parseResult == PASSEDchangeRoleToLeader。注意它把 REJECT_TERM_NOT_READY(对方 term 还没到位)也计入 acceptedNum + notReadyTermNum 的多数派,允许「等一轮再投」而不是直接失败,这让选举在 term 不一致时也能快速收敛。

总结

DLedger 用 Raft 的多数派提交替代异步复制,将”主库故障切换”从手动操作变为自动选举。代价是每条消息的写入延迟增加(需等多数派 ACK),吞吐约下降 10-15%。

章末提问

Q1:DLedger 为什么要用 Raft 替代异步主从复制? 答:因为异步复制下「已 ACK」不等于「已持久化到多数节点」。因为主库本地落盘就回 ACK,主库宕机后未同步到从库的已 ACK 消息永久丢失,且故障切换要靠人工。Raft 多数派提交把「提交」绑定到「多数派落盘」,做到已 ACK 不丢、故障自动选举。

Q2:多数派提交怎么保证已 ACK 的消息不丢? 答:因为只有多数派节点落盘后 Leader 才 commit 并回 ACK,所以这条已 ACK 的消息至少存在于一个「多数派」集合里;任意少数节点宕机后,存活节点中必然包含该多数派里的至少一个节点。因为 Raft 保证新 Leader 拥有所有已 committed 日志,所以切主后这条消息依然存在,已 ACK 不丢。

Q3:为什么 ConsumeQueue 和 IndexFile 不参与 Raft 同步? 答:因为它们是 CommitLog 的派生索引,可以从 CommitLog 完整重建,不是源数据。因为纳入 Raft 会让共识数据量膨胀 3-5 倍、拖慢同步,所以只同步 CommitLog,索引由各节点本地异步重建,Follower 晋升 Leader 后等索引重建完成再对外提供消费服务。

Q4:选主为什么是 lastLogIndex 最大的节点当选? 答:因为日志最旧的节点若当选 Leader,会把落后数据当作权威,导致已 committed 的消息被覆盖丢失。因为 Raft 规定「日志最新者才有资格当选」,所以新 Leader 一定包含所有已 committed 日志,切主不丢已提交数据;handleVote 里先比 ledgerEndTerm、再比 ledgerEndIndex 就是实现这条约束。

Q5:DLedger 的代价是什么? 答:代价是每条消息写延迟增加、吞吐约降 10-15%。因为每条消息都要等多数派 Follower 落盘 ACK 才能返回,网络往返和落盘等待时间计入延迟,且 Leader 要并行维护向所有 Follower 的复制。


Share this post on:

Previous Post
Redis Cluster Gossip协议——去中心化的元数据传播
Next Post
双亲委派机制——为什么你不能自定义java.lang.String