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 的只有一层集成代码DLedgerCommitLog(org.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.isQuorum 和 DLedgerLeaderElector.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 == PASSED 就 changeRoleToLeader。注意它把 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 的复制。