RocketMQ 面试回答
总纲(30s 开场白)
RocketMQ 是一个面向在线业务消息的分布式消息中间件。它的核心设计决策是「所有 Topic 的消息共用一个 CommitLog 顺序追加写」——因为顺序写比随机写快百倍,用 ConsumeQueue 二级索引把读性能补回来;再靠「NameServer 无状态路由 + 主从/DLedger 复制」保证高可用。一句话:用存储结构的取舍换极致的写吞吐,用复制冗余换消息不丢。
① 架构设计(NameServer / Broker / Producer / Consumer)
一句话结论(30s)
RocketMQ 是「路由与存储解耦」的四角色架构:NameServer 只做无状态路由发现、Broker 只做消息存储投递、Producer/Consumer 只做收发。最关键的设计决策是所有 Topic 的消息写进同一个 CommitLog 顺序追加——因为磁盘顺序写比随机写快 100 倍以上,代价是读消息需要 ConsumeQueue 二级索引做一次定位。
核心原理(2min)
四个角色各司其职:
| 角色 | 职责 | 关键点 |
|---|---|---|
| NameServer | 路由注册中心 | 无状态、各节点对等互不通信,维护 Topic→Broker 路由表 |
| Broker | 消息存储 + 投递 | 三层存储(CommitLog/ConsumeQueue/IndexFile)、主从复制、刷盘、存消费进度 |
| Producer | 生产消息 | 从 NameServer 拉路由,按队列选择策略发消息,失败自动重试 |
| Consumer | 消费消息 | 从 NameServer 拉路由,长轮询拉消息,消费进度提交给 Broker |
- NameServer 为什么无主、不用 ZooKeeper? 先自己想:注册中心到底需不需要强一致?答案是「路由信息允许秒级不一致」——Broker 每 30s 上报一次心跳,NameServer 每 10s 扫描、超过 120s 没收到心跳就下线该 Broker。顺着推下去:既然允许短暂不一致,那「选主」这种强一致才需要的重活就纯属多余,于是用「允许短暂不一致」换「零运维」,NameServer 可随时重启、动态扩缩,不需要选主。
- 数据流:Producer 拿到路由 → 选一个 MessageQueue 发消息 → Broker 写入 CommitLog 并构建 ConsumeQueue → Consumer 长轮询从 Broker 拉消息 → 处理完提交 offset(进度存在 Broker 上)。
底层深入(5-10min)
取舍一:CommitLog 单文件顺序追加。 先想:100 个 Topic 的消息都堆进同一个文件,看起来更乱,凭什么反而更快?——因为所有 Topic、所有队列的消息统一追加到同一个 CommitLog 文件末尾(满 1GB 新建一个文件)。如果每个 Topic 一个文件,100 个 Topic 就是 100 处随机写,磁头来回寻道;单文件顺序写 = 1 次顺序 IO,HDD 顺序写比随机写快约 4 倍,SSD 下也省了 IO 调度开销。
真实源码里,写入路径是「取最后一个 MappedFile → 加 putMessageLock → appendMessage」,CommitLog.asyncPutMessage 关键片段(store/CommitLog.java):
MappedFile mappedFile = this.mappedFileQueue.getLastMappedFile();
long currOffset;
if (mappedFile == null) {
currOffset = 0;
} else {
currOffset = mappedFile.getFileFromOffset() + mappedFile.getWrotePosition();
}
...
putMessageLock.lock(); //spin or ReentrantLock, depending on store config
try {
...
if (null == mappedFile || mappedFile.isFull()) {
mappedFile = this.mappedFileQueue.getLastMappedFile(0); // Mark: NewFile may be cause noise
...
}
...
result = mappedFile.appendMessage(msg, this.appendMessageCallback, putMessageContext);
switch (result.getStatus()) {
case PUT_OK:
onCommitLogAppend(msg, result, mappedFile);
break;
case END_OF_FILE:
onCommitLogAppend(msg, result, mappedFile);
unlockMappedFile = mappedFile;
// Create a new file, re-write the message
mappedFile = this.mappedFileQueue.getLastMappedFile(0);
...
result = mappedFile.appendMessage(msg, this.appendMessageCallback, putMessageContext);
...
break;
...
}
...
} finally {
...
putMessageLock.unlock();
}
getLastMappedFile() 永远取队列里最后一个文件,写满(isFull())或触发 END_OF_FILE 时立刻换下一个文件重写——这就是「单文件顺序追加」在代码层面的落点。加锁(putMessageLock)是为了串行化并发写,保证 offset 严格连续递增,不产生空洞。
取舍二:ConsumeQueue 二级索引(20 字节的精确定位)。 写是顺序了,那读怎么办?先想:消息全混在一个 CommitLog 里,消费者总不可能顺序扫描去捞自己的消息吧?——于是用 ConsumeQueue 稀疏索引做定位,每条固定 20 字节 = 8B CommitLog Offset + 4B Size + 8B Tag Hash:
CommitLog: [Msg_A_Q0][Msg_A_Q1][Msg_B_Q0][Msg_A_Q2]... ← 按到达顺序
ConsumeQueue(TopicA/Queue0): [offset,size,tagHash][offset,size,tagHash]... ← 每 Queue 一个文件
读取链路:消费者先顺序读 ConsumeQueue(快)→ 拿到 offset/size → 去 CommitLog 用 FileChannel.position(offset) 做一次随机读定位完整消息。即「写全顺序、读仅一次随机」。Tag 过滤也在 ConsumeQueue 层完成(用 8B 的 hashCode 匹配),不匹配的不用读 CommitLog,省 IO。
ConsumeQueue 是异步构建的(后台 ReputMessageService 每 1ms 检测 CommitLog 新消息后追加索引),也是可重建的——Broker 启动时若发现非正常关闭(abort 文件存在),会从 CommitLog 重扫重建所有 ConsumeQueue 和 Index。先想:它为什么敢异步构建、敢「可丢」?——因为索引只是 CommitLog 的影子,原始数据还在。所以它是最「可丢」的一层,丢了能重建。
20 字节这个数是写死在源码里的常量,ConsumeQueue.java 里直接注明「Offset(8) + Size(4) + TagHashCode(8) = 20 Bytes」,写入时也严格按这个顺序 put 三个字段:
/**
* ConsumeQueue's store unit. Size: CommitLog Physical Offset(8) + Body Size(4) + Tag HashCode(8) = 20 Bytes
*/
public static final int CQ_STORE_UNIT_SIZE = 20;
...
this.byteBufferIndex.flip();
this.byteBufferIndex.limit(CQ_STORE_UNIT_SIZE);
this.byteBufferIndex.putLong(offset);
this.byteBufferIndex.putInt(size);
this.byteBufferIndex.putLong(tagsCode);
final long expectLogicOffset = cqOffset * CQ_STORE_UNIT_SIZE;
byteBufferIndex 是 ByteBuffer.allocate(20) 的复用缓冲区,putLong/putInt/putLong 三个字段合计正好 20 字节、逐条追加进 ConsumeQueue 文件。expectLogicOffset = cqOffset * 20 说明每条索引在文件里偏移严格对齐,读时能按 i += 20 直接步进定位,无需额外元数据。
取舍三:NameServer 无主设计。 各节点互不通信,每个节点都有全量路由;Broker 向所有 NameServer 心跳上报。这里再想一次:选主是为了强一致,可前面已经论证了「路由允许秒级不一致」,那选主就纯属多余。对比 ZooKeeper 的 Raft 选主(强一致 + 高运维成本),NameServer 用「允许秒级不一致」换「零运维」,这是「够用就好」的典型工程取舍。
② 消息可靠(不丢 / 刷盘 / 事务消息)
一句话结论(30s)
RocketMQ 保证消息不丢靠三段接力:Producer 端「同步发送 + 失败重试」、Broker 端「同步刷盘 + 主从复制」、Consumer 端「消费成功才提交 offset + 失败重试」。可靠性的本质是用「多等一次磁盘 IO / 多等一个副本 ACK」换「不丢」。
核心原理(2min)
先想:一条消息从生产者到消费者,会在哪几个环节「可能丢」?顺着链路数一遍——发送出去没确认、Broker 没落盘就宕机、消费没成功就推进度,正好三段接力:
- Producer 端:
send()返回SendResult才算发送成功,否则自动重试(默认重试 2 次);失败也能拿到本地重试。 - Broker 端(刷盘):
- 同步刷盘 SYNC_FLUSH:消息进 Page Cache 后调用
fsync()强制落盘,落盘成功才返回 ACK。可靠但吞吐低(几百到几千 TPS)。 - 异步刷盘 ASYNC_FLUSH(默认):进 Page Cache 就返回 ACK,后台线程 500ms 批量刷盘。吞吐高(数万到十万 TPS),但 OS 崩溃会丢 Page Cache 里的消息。
- 同步刷盘 SYNC_FLUSH:消息进 Page Cache 后调用
- Consumer 端:默认集群模式消费成功后才提交 offset;处理失败返回
RECONSUME_LATER,消息进重试队列,不丢、等重试。 - 事务消息:用 Half Message 两阶段提交,保证「本地事务成功才投递消息」,解决「扣库存 + 发消息」的原子性问题。
底层深入(5-10min)
为什么顺序写能到数十万 TPS? 先想:同步刷盘明明每次都 fsync()、一次磁盘 IO 就是几毫秒,凭什么还能跑到几十万 TPS?——答案藏在三个关键词里:mmap + 顺序写 + GroupCommit。
- mmap 内存映射文件:CommitLog 用
MappedFile映射到虚拟地址空间,用户态直接写 Page Cache,绕过用户态 Buffer(传统write()两次拷贝,mmap 一次)。 - 顺序写:Append-only 让磁头不动,SSD 顺序写 600MB/s、HDD 约 100MB/s,而随机写(如 InnoDB B+Tree 更新)只有几 MB/s,差 100 倍以上。
- GroupCommit 组提交:同步刷盘每次写都
fsync()太慢,把 10ms 内汇聚的一批写请求合并成一次 fsync,一次刷多条,把同步刷盘吞吐拉高近 10 倍。看明白这层就懂了:不是「少刷」,而是「攒批一次刷多条」。
mmap 映射和落盘都是 DefaultMappedFile 干的:初始化时 FileChannel.map 把文件映射到虚拟地址空间,刷盘时用 force() 把 Page Cache 强制写回磁盘:
// DefaultMappedFile.init —— mmap 映射
this.fileChannel = new RandomAccessFile(this.file, "rw").getChannel();
if (writeWithoutMmap) {
// Still create MappedByteBuffer for reading operations
this.mappedByteBuffer = this.fileChannel.map(MapMode.READ_ONLY, 0, fileSize);
} else {
// Use MappedByteBuffer for both reading and writing (default behavior)
this.mappedByteBuffer = this.fileChannel.map(MapMode.READ_WRITE, 0, fileSize);
}
// DefaultMappedFile.flush —— 落盘
public int flush(final int flushLeastPages) {
if (!isWriteable()) {
return this.getFlushedPosition();
}
if (this.isAbleToFlush(flushLeastPages)) {
if (this.hold()) {
int value = getReadPosition();
try {
this.mappedByteBufferAccessCountSinceLastSwap++;
//We only append data to fileChannel or mappedByteBuffer, never both.
if (writeWithoutMmap || writeBuffer != null || this.fileChannel.position() != 0) {
this.fileChannel.force(false);
} else {
this.mappedByteBuffer.force();
}
this.lastFlushTime = System.currentTimeMillis();
FLUSHED_POSITION_UPDATER.set(this, value);
...
fileChannel.map(MapMode.READ_WRITE, ...) 返回的 MappedByteBuffer 让用户态直接写 Page Cache,省掉 write() 的一次拷贝;force() 等价于 fsync,把脏页刷到磁盘并推进 FLUSHED_POSITION。异步刷盘就是「后台线程定时调 flush」,同步刷盘则是在 ACK 前先调一次 flush——二者共用同一份代码,只是触发时机不同。
事务消息(Half Message 两阶段提交):先想:扣库存(DB)和发消息(MQ)是两套系统,怎么保证「要么都成功、要么都不做」?——顺序调用两个系统做不到原子,于是引入一个「消费者看不到的中间态」Half Message:
1. 发送 Half Message(写进内部 Topic RMQ_SYS_TRANS_HALF_TOPIC,对 Consumer 不可见)
2. 执行本地事务(如扣库存)
3. 成功 → Commit(转写目标 Topic,Consumer 可见);失败 → Rollback(不投递)
- Half Message 不写目标 Topic,写系统内置
RMQ_SYS_TRANS_HALF_TOPIC,消息属性里记原 Topic。
Half Message 的「换 Topic + 记原 Topic」在 TransactionalMessageBridge.parseHalfMessageInner 里一行行落地(broker/transaction/queue/TransactionalMessageBridge.java):
private MessageExtBrokerInner parseHalfMessageInner(MessageExtBrokerInner msgInner) {
String uniqId = msgInner.getUserProperty(MessageConst.PROPERTY_UNIQ_CLIENT_MESSAGE_ID_KEYIDX);
if (uniqId != null && !uniqId.isEmpty()) {
MessageAccessor.putProperty(msgInner, TransactionalMessageUtil.TRANSACTION_ID, uniqId);
}
MessageAccessor.putProperty(msgInner, MessageConst.PROPERTY_REAL_TOPIC, msgInner.getTopic());
MessageAccessor.putProperty(msgInner, MessageConst.PROPERTY_REAL_QUEUE_ID,
String.valueOf(msgInner.getQueueId()));
msgInner.setSysFlag(
MessageSysFlag.resetTransactionValue(msgInner.getSysFlag(), MessageSysFlag.TRANSACTION_NOT_TYPE));
...
msgInner.setTopic(TransactionalMessageUtil.buildHalfTopic());
msgInner.setQueueId(0);
msgInner.setPropertiesString(MessageDecoder.messageProperties2String(msgInner.getProperties()));
return msgInner;
}
先把真实 Topic 存进 PROPERTY_REAL_TOPIC 属性,再把消息的 Topic 改成 buildHalfTopic()(即 RMQ_SYS_TRANS_HALF_TOPIC)并强制 queueId=0——消费者订阅的是真实 Topic,永远拉不到半消息。Commit 时由 EndTransactionProcessor 调 MessageAccessor.clearProperty(msgInner, PROPERTY_TRANSACTION_PREPARED) 清除「准备中」标记后 sendFinalMessage 转投真实 Topic,才对外可见。
- 回查机制兜底网络丢包:先想:如果 Commit 这步网络超时、Broker 没收到,Half Message 会不会永远悬空?——不会,
TransactionalMessageCheckService每 60s 扫描,找到超过transactionTimeOut(默认 6s)还未确认的 Half Message,向原 Producer 发 CHECK,调checkLocalTransaction()返回 COMMIT/ROLLBACK/UNKNOWN;最多回查 15 次,超了自动 Rollback,防止永久悬空。 - 已处理的消息 offset 记在
RMQ_SYS_TRANS_OP_HALF_TOPIC,回查时先查它避免重复投递。 - 即使事务消息保证「先提交本地事务才投递」,消费端仍可能重复(重试、回查),所以消费端必须幂等——这正好引出下面的幂等主题。
③ 集群高可用(主从 / DLedger)
一句话结论(30s)
RocketMQ 高可用分两代:传统「主从异步复制」靠主从 + 手动切换,DLedger 把 Raft 共识嵌进 CommitLog 写入路径、用「多数派提交」实现自动故障切换。核心是用「多等一个副本的 ACK」换「主库挂了消息也不丢、能自动选主」。
核心原理(2min)
先想:都做了主从复制了,为什么消息还会丢、还非要人工切换?——问题出在「异步」:Master 写完立即 ACK,此时数据可能还没同步到 Slave,Master 一宕,已 ACK 但未同步的就丢了,Slave 也说不清自己是不是最新、不敢自动上位。
- 主从复制:Master 写 CommitLog 后把数据复制给 Slave,Slave 只读、可分担消费读压力。传统模式是异步复制——Master 写完立即 ACK,异步同步给 Slave,Master 宕机时已 ACK 但未同步的消息会丢,且需人工提升 Slave。
- DLedger:将 Raft 共识协议嵌入写入路径,3 节点中多数派(≥2)写入成功才 commit 并返回 Producer。Leader 宕机后 Follower 选举超时自动发起
RequestVote,lastLogIndex最大的节点当选,对 Producer 透明。 - 消费进度存在 Broker 上随主从同步,切换后进度不丢。
底层深入(5-10min)
DLedger 五步写入:
Producer → Leader Broker:
1. Leader 写本地 DLedger Log(即 CommitLog Entry)
2. Leader 并行向 Follower1/Follower2 发 AppendEntries RPC
3. Follower 写本地日志 → 返回 ACK
4. Leader 收到多数派(含自身)ACK → commit(更新 commitIndex)
5. Leader 返回 Producer 写入成功
DLedgerCommitLog.asyncPutMessage 里「写本地 + 等多数派」的代码是这么接起来的(store/dledger/DLedgerCommitLog.java):
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, ...));
}
long wroteOffset = dledgerFuture.getPos() + DLedgerEntry.BODY_OFFSET;
...
return dledgerFuture.thenApply(appendEntryResponse -> {
...
switch (DLedgerResponseCode.valueOf(appendEntryResponse.getCode())) {
case SUCCESS:
putMessageStatus = PutMessageStatus.PUT_OK;
break;
...
case WAIT_QUORUM_ACK_TIMEOUT:
putMessageStatus = PutMessageStatus.IN_SYNC_REPLICAS_NOT_ENOUGH;
break;
...
handleAppend(request) 一个调用把「写本地日志 + 复制给 Follower + 等多数派 ACK」全包了,返回的 Future 只有在多数派确认后才 resolve。注意 WAIT_QUORUM_ACK_TIMEOUT(等多数派超时)被映射成 IN_SYNC_REPLICAS_NOT_ENOUGH 返回给 Producer——这就是「多数派没凑齐就不算写入成功」在返回码上的体现。
为什么 ConsumeQueue 不参与 Raft? 先想:Raft 要复制的是「原始数据」,还是「能从原始数据重建出来的东西」?——ConsumeQueue 和 IndexFile 是 CommitLog 的派生索引,能从 CommitLog 完全重建,纳入 Raft 会让共识数据量膨胀 3-5 倍。所以 DLedger 只同步 CommitLog,ConsumeQueue/IndexFile 由每个 Broker 本地异步重建——Follower 晋升为 Leader 后,要等索引重建完成才能对外提供消费服务。
选主与容错:Raft 保证新 Leader 一定拥有所有已 committed 的日志(lastLogIndex 最大者当选);原 Leader 恢复后自动降级为 Follower。3 节点容忍 1 故障,5 节点容忍 2 故障。先想:白拿「自动选主 + 消息不丢」有没有代价?——当然有,代价是每条消息要等多数派 ACK,延迟增加,吞吐约下降 10-15%——这是「可用性换一致性」的典型 trade-off。
④ 死信队列
一句话结论(30s)
死信队列是 RocketMQ 的「最后防线」:消息重试 16 次仍失败就转入
%DLQ%{ConsumerGroup}、永久不再重试。核心设计是用「延迟递增重试 + 16 次上限」既防止毒消息无限重试浪费资源,又把问题消息隔离出正常消息流。
核心原理(2min)
先想:一条消息怎么处理都失败(比如业务有 bug),你会无限重试还是直接丢弃?——无限重试会拖垮队列、浪费资源,直接丢弃又丢数据,所以 RocketMQ 走中间路:延迟递增重试 + 16 次上限 + 隔离。
- 消费失败返回
RECONSUME_LATER→ 消息进%RETRY%{ConsumerGroup}重试 Topic,延迟递增重试:1s、5s、10s、30s … 直到第 16 次约 2h。 - 16 次全部失败 → 转入
%DLQ%{ConsumerGroup},永久不重试,隔离出正常流。 - 标准处理流程四步:告警(订阅 DLQ 触发钉钉/企微)→ 排查(业务日志 + 消息内容)→ 修复(修 bug 或补偿下游数据)→ 回放(
admin.resendMessageByHandleOffset()把 DLQ 消息重新投递到原始 Topic)。
底层深入(5-10min)
为什么延迟递增而不是等间隔重试? 消费失败的根因(bug、下游宕机、资源不足)不会短时间自愈,等间隔重试会在下游恢复前反复失败、白白浪费资源。延迟递增让重试逐步退避:前几次间隔短(应对网络抖动这种瞬时故障,快速重试),后几次间隔长(下游长时间宕机时给足恢复时间);16 次上限保证不会无限重试阻塞队列。
「16 次上限 → 转 DLQ」的判定不在 Consumer 端,而在 Broker 端 SendMessageProcessor.handleRetryAndDLQ(broker/processor/SendMessageProcessor.java):
int maxReconsumeTimes = subscriptionGroupConfig.getRetryMaxTimes();
...
int reconsumeTimes = requestHeader.getReconsumeTimes() == null ? 0 : requestHeader.getReconsumeTimes();
...
if (reconsumeTimes > maxReconsumeTimes || sendRetryMessageToDeadLetterQueueDirectly) {
...
properties.put(MessageConst.PROPERTY_DELAY_TIME_LEVEL, "-1");
newTopic = MixAll.getDLQTopic(groupName);
int queueIdInt = randomQueueId(DLQ_NUMS_PER_GROUP);
topicConfig = this.brokerController.getTopicConfigManager().createTopicInSendMessageBackMethod(newTopic,
DLQ_NUMS_PER_GROUP,
PermName.PERM_WRITE | PermName.PERM_READ, 0
);
msg.setTopic(newTopic);
msg.setQueueId(queueIdInt);
msg.setDelayTimeLevel(0);
...
}
Consumer 消费失败 sendMessageBack 时会把消息发回 %RETRY%{Group},请求头带上 reconsumeTimes;Broker 拿它和 maxReconsumeTimes(订阅组默认 16)比,一旦超限就把消息 Topic 改成 MixAll.getDLQTopic(groupName)(即 %DLQ%{Group})并 DELAY_TIME_LEVEL=-1 取消延时。Topic 命名就在 MixAll 里两条常量拼接:RETRY_GROUP_TOPIC_PREFIX + consumerGroup、DLQ_GROUP_TOPIC_PREFIX + consumerGroup,%RETRY% 和 %DLQ% 都是普通 Topic,复用同一套存储。
几个要点:
- 重试 Topic 本质是普通 Topic,重试消息会重新写进
%RETRY%对应的 ConsumeQueue,按延迟级别定时投递——它复用了延时消息的机制(同一套 CommitLog + ConsumeQueue 存储引擎)。 - DLQ 消息默认保留 3 天(
fileReservedTime=72h,可配),给运维足够的排查窗口。 - 不要在原 ConsumerGroup 直接重消费 DLQ——先想:bug 没修就把消息塞回原队列会怎样?——又失败、又进一次 DLQ,形成死循环。所以先把 bug 修好再手动回放。
⑤ 顺序消息 / 幂等(引申)
顺序消息
一句话结论(30s):RocketMQ 顺序消息靠三环接力——生产者取模路由到同一队列、Broker 队列 FIFO 存储、Consumer 用 MessageListenerOrderly 加锁 + 单线程串行消费——本质是把「必须有序」的消息串行化到最小范围。
核心原理(2min):
- 先分清「有序」两个层次:全局有序(整个 Topic 一个队列,完全丧失并行度,很少用)和局部有序(同一业务标识如订单 ID 的消息有序,不同标识并行——生产主流)。
- 三步:先想:要让「同一订单的消息严格按顺序消费」,靠哪一环单打独斗能成吗?——不能,必须生产、存储、消费三环接力,把消息串行化到最小范围:① 生产者用
MessageQueueSelector按orderId % queueSize(或hashCode % size)路由,保证同一订单永远进同一队列;② Broker 每个 MessageQueue 是 FIFO,顺序写天然保证「发送顺序 = 存储顺序 = 拉取顺序」;③ Consumer 用MessageListenerOrderly:单队列加分布式锁(同一时刻一个 Queue 只能被一个实例的一个线程消费)+ 同步逐条消费(前一条 SUCCESS 才处理下一条)。
底层深入(5-10min):
- Rebalance 是顺序消费最大的敌人:先想:既然一个 Queue 一把锁、只能一个线程消费,那消费者扩缩容要重新分配队列时,锁怎么安全交接?——Consumer 扩缩容触发重平衡时,新 Consumer 要先向 Broker 抢队列锁,旧 Consumer 不释放锁新 Consumer 就消费不了,从而保证「任意时刻一个 Queue 最多一个线程消费」;旧 Consumer 挂了则靠锁超时(默认 30s)自动释放。代价是 Rebalance 期间该队列暂停几百毫秒。
SUSPEND_CURRENT_QUEUE_A_MOMENT:顺序消费失败不能跳过这条(否则乱序),只能挂起整个队列、setSuspendCurrentQueueTimeMillis(3000)后重试同一批;达到 16 次上限进死信。代价:一条消息卡住,整个队列后面的全卡住——所以每条处理逻辑必须设超时并告警。
「挂起重试而非跳过」在消费端 ConsumeMessageOrderlyService.processConsumeResult 里实现(client/impl/consumer/ConsumeMessageOrderlyService.java):
case SUSPEND_CURRENT_QUEUE_A_MOMENT:
this.getConsumerStatsManager().incConsumeFailedTPS(consumerGroup, consumeRequest.getMessageQueue().getTopic(), msgs.size());
if (checkReconsumeTimes(msgs)) {
consumeRequest.getProcessQueue().makeMessageToConsumeAgain(msgs);
this.submitConsumeRequestLater(
consumeRequest.getProcessQueue(),
consumeRequest.getMessageQueue(),
context.getSuspendCurrentQueueTimeMillis());
continueConsume = false;
} else {
commitOffset = consumeRequest.getProcessQueue().commit();
}
break;
continueConsume = false 是关键——顺序消费失败后不 commit、不推进 offset,而是 submitConsumeRequestLater(..., suspendCurrentQueueTimeMillis) 把整个队列延后重投,下一批还从同一条开始。checkReconsumeTimes(msgs) 为 false(达到重试上限)时走 commit() 直接提交并转 sendMessageBack 进死信,而不是卡死。
- 并行度 = 队列数:Queue 太少吞吐上不去,太多浪费句柄;一般设为「Consumer 节点数 × 核数的 2-4 倍」。核心心法:先想清楚「顺序」和「并行」是不是天生的死对头?——是,顺序 = 串行化,要顺序就没有完全并行,把需要顺序的放一个 Queue、不需要的放其他 Queue。
消息幂等
一句话结论(30s):RocketMQ 投递语义是 at-least-once(至少一次),重复是设计如此、不是 bug;所以幂等是消费端必须自己兜底的最后防线——用业务唯一标识 + DB 唯一索引 / Redis SETNX / 业务状态机去重,不能依赖 MessageId。
核心原理(2min):先想:消息为什么好端端会重复?——本质是链路上任何一个「确认 ACK」丢了,对方就以为没成功、重来一遍。三个重复来源——
- 生产者重试发送:ACK 丢失,Producer 不知道 Broker 是否已写入,重发导致同一业务数据两份(MessageId 不同)。
- 消费者 ACK 丢失:提交 offset 时网络超时,Broker 以为没消费成功,重新投递。
- Rebalance:消费到一半队列被重新分配,新 Consumer 从头拉,重复消费。
底层深入(5-10min):
- 为什么 MessageId 不可靠:先想:MessageId 每条消息都有、看起来天然能当去重键,用它不是最省事吗?——不行,MessageId 是 Broker 写入 CommitLog 时生成的(含 Broker 地址 + offset),生产者重试发送会得到不同 MessageId,所以它只能识别「同一条消息被消费两次」,识别不了「同一业务数据被发两次」。去重必须用业务唯一键(订单号、流水号、traceId)。
- 四种去重方案:
- DB 唯一索引(最可靠):去重表以 biz_id 为主键,先
INSERT,DuplicateKeyException说明已处理直接返回 SUCCESS;与业务库共享事务、强一致。 - Redis SETNX(高性能):
set(dedupKey, "1", "NX", "EX", 3600),key 已存在即重复;缺点是 Redis 与业务库不在同一事务。 - 业务状态机(最优雅):
UPDATE orders SET status='PAID' WHERE order_id=? AND status='PENDING',affected==0说明已处理过,零额外成本。 - RocketMQ 5.0 Broker 端去重:用
messageKey去重,但 Broker 有额外开销,尚未大规模生产验证。
- DB 唯一索引(最可靠):去重表以 biz_id 为主键,先
- 去重记录要定期清理(DB 表定时删 3 天前的,Redis 设 TTL 72h,覆盖所有重试场景),避免无限增长。
追问清单(10 问:一句话结论 + 展开)
1. NameServer、Broker、Producer、Consumer 分别干什么?
一句话结论:NameServer 管路由、Broker 管存储投递、Producer 生产、Consumer 消费,四者通过「路由发现 + 直连 Broker」解耦。
NameServer 是无状态对等节点,维护 Topic→Broker 路由表,Broker 每 30s 心跳上报、NameServer 120s 没收到就下线。Producer/Consumer 启动时和定时从 NameServer 拉路由并本地缓存,之后直连 Broker 收发消息,所以 NameServer 不参与数据链路,只做元数据。Broker 负责三层存储(CommitLog/ConsumeQueue/IndexFile)、刷盘、主从复制,还存消费进度。
2. RocketMQ 和 Kafka 的核心区别?为什么选 RocketMQ?
一句话结论:Kafka 为「高吞吐流式日志」而生、怕多 Topic;RocketMQ 为「在线业务消息的可靠性」而生、多 Topic 不衰减——选 RocketMQ 因为需要事务/延时消息、Tag 过滤和消息不丢。
存储模型是根因:Kafka 每个 Partition 一个独立目录,1000 个 Topic 就是几千个文件句柄 + 随机 IO,吞吐可能掉 50%;RocketMQ 所有 Topic 共用一个 CommitLog 顺序写,Topic 增多不影响写性能,靠 ConsumeQueue 索引定位。功能上 RocketMQ 原生支持事务消息、延时消息(18 级)、Tag/SQL 过滤、死信队列、消息重试,Kafka 要么没有要么需额外开发。所以「一条消息丢了会不会半夜叫醒 tech lead?会就选 RocketMQ」。
3. 怎么保证消息不丢?生产者、Broker、消费者各做了什么?
一句话结论:三段接力——Producer 同步发送 + 重试、Broker 同步刷盘 + 主从复制、Consumer 消费成功才提交 offset。
Producer 端 send() 拿到 SendResult 才确认成功,失败自动重试(默认 2 次)。Broker 端同步刷盘 fsync() 落盘才 ACK,再叠加主从复制让副本兜底(即使 Master 宕机 Page Cache 丢,Slave 还有)。Consumer 端集群模式先消费成功、再提交 offset,失败返回 RECONSUME_LATER 进重试队列,不确认就不丢。三段各守一段,任一环节用默认配置都可能丢(如默认异步刷盘 + 异步复制)。
4. 同步刷盘和异步刷盘的区别?性能差多少?
一句话结论:同步刷盘 =
fsync()落盘后才 ACK(不丢但慢,几百几千 TPS),异步刷盘 = 进 Page Cache 就 ACK(快但 OS 崩溃可能丢,数万十万 TPS),差一到两个数量级。
同步刷盘每次写都要等一次磁盘 IO,延迟高、吞吐低,靠 GroupCommit(10ms 汇聚一批合并成一次 fsync)把吞吐拉高近 10 倍;异步刷盘是默认,后台线程 500ms 批量刷盘,Broker 进程崩溃一般不丢(OS 会继续刷),只有 OS 崩溃才丢。生产折中:多数用「异步刷盘 + 同步主从复制」,既保住吞吐又有副本兜底;金融支付才上同步刷盘。
5. Broker 主从怎么切换?DLedger 解决什么问题?
一句话结论:传统模式主从是异步复制 + 人工切换,DLedger 用 Raft 多数派提交实现自动选举,解决「主库挂了消息不丢 + 免人工」两个问题。
传统模式 Master 写完立即 ACK 再异步同步 Slave,Master 宕机时已 ACK 未同步的消息永久丢,且要手动提升 Slave。DLedger 把 Raft 嵌进写入路径:Leader 并行发 AppendEntries,多数派(3 节点 ≥2)写入成功才 commit 并返回 Producer;Leader 宕机后 Follower 选举超时发起 RequestVote,lastLogIndex 最大者当选,对 Producer 透明。代价是吞吐下降 10-15%。
6. 如果 NameServer 全挂了,已有的消费者还能消费吗?
一句话结论:能。因为路由是本地缓存 + 直连 Broker,消费过程不经过 NameServer。
Consumer 从 NameServer 拉到路由后会缓存到本地,并和 Broker 建立长连接;拉消息时直接用缓存路由 + 已建立的连接长轮询,NameServer 不参与。NameServer 全挂只影响:新 Producer/Consumer 起不来(拿不到路由)、Broker 上下线/扩缩容等路由变化无法感知(Rebalance 卡住)。所以「已有消费者继续消费、但别动集群拓扑」是正确认知。
7. 什么情况消息进死信?怎么处理死信?
一句话结论:消费失败返回
RECONSUME_LATER后按延迟递增重试,16 次仍失败就进%DLQ%{Group}永久不再重试。
重试间隔 1s→5s→10s→30s…→2h 逐步退避,16 次上限防毒消息无限重试。处理死信四步:订阅 DLQ 触发告警 → 用业务日志 + 消息体排查根因 → 修 bug 或补偿下游 → resendMessageByHandleOffset() 手动回放原始 Topic。
8. 顺序消息怎么实现?哪些场景需要顺序消息?
一句话结论:三环接力——生产者
MessageQueueSelector按订单 ID 取模路由到同一队列、Broker 队列 FIFO、ConsumerMessageListenerOrderly加锁 + 单线程串行。
生产按 orderId % queueSize(注意自增 ID 要 hashCode 再取模避免热点)保证同一订单进同一队列;Broker 顺序写天然 FIFO;消费端单队列加分布式锁 + 逐条 SUCCESS 才下一条,失败 SUSPEND_CURRENT_QUEUE_A_MOMENT 挂起重试而非跳过。场景:订单的「创建→支付→发货→完成」状态流转、数据库 binlog 同步,都需要同一业务标识严格有序。
9. 消息重复消费怎么处理?怎么保证幂等?
一句话结论:RocketMQ 是 at-least-once,重复来自「发送重试 / ACK 丢失 / Rebalance」三处,消费端用业务唯一键 + DB 唯一索引或 Redis SETNX 或状态机兜底。
MessageId 不能去重(生产者重试会产生不同 MessageId),必须带订单号/流水号等业务键。最可靠是 DB 去重表 INSERT 撞 DuplicateKeyException 判重;最高性能是 Redis SET NX EX;最优雅是状态机 UPDATE ... WHERE status='PENDING' 的 affected==0 判重。实践上常用「业务唯一键唯一索引 + 业务状态机」双保险,去重记录定时清理。
10. 事务消息原理?什么场景会用到?
一句话结论:Half Message 两阶段提交——先发不可见的 Half Message、再执行本地事务、最后 Commit/Rollback,配「回查机制」兜底网络丢包。
Half Message 写内部 Topic RMQ_SYS_TRANS_HALF_TOPIC(Consumer 不可见),本地事务成功后 Commit 转写真实 Topic、失败 Rollback。若 Commit/Rollback 因网络超时没到 Broker,TransactionalMessageCheckService 每 60s 扫描、超过 6s 未确认的 Half Message 回查 Producer 的 checkLocalTransaction(),最多回查 15 次、超了自动 Rollback。典型场景:下单 → 扣库存(本地事务)→ 发订单消息,需要 DB 和 MQ 的原子性;若未直接使用事务消息,也可用本地消息表 + 定时补偿实现等价效果。