RocketMQ 集群运维

共 20 题
📑 题目列表 20 题
#
★★★

1. Broker 主从同步中 SYNC_MASTER/ASYNC_MASTER、同步双写与异步复制的数据安全边界

RocketMQ 的默认主从同步模式(SYNC_MASTER/ASYNC_MASTER)如何决定数据安全边界,同步双写与异步复制有何区别?

  • SYNC_MASTER 与 ASYNC_MASTER 的区别
  • 同步双写与异步复制的数据安全边界
  • 主从切换时数据丢失可能性

RocketMQ 的 Broker 复制方式由 brokerRole 决定:SYNC_MASTER 表示主从同步复制,Master 收到消息后需由 Slave 确认写入成功才返回给生产者;ASYNC_MASTER 表示主从异步复制,Master 收到消息即返回,Slave 异步复制,主节点宕机时 Slave 可能丢失尚未复制的最新消息。数据安全边界取决于"同步复制 + 同步刷盘"的组合:SYNC_MASTER + SYNC_FLUSH 在 Master 与 Slave 都刷盘后才返回,可靠性最高;而 ASYNC_MASTER + ASYNC_FLUSH 在 Master 宕机时可能丢数据。运维上,对账务等强一致场景应选择 SYNC_MASTER,对吞吐优先场景可选 ASYNC_MASTER。

同步复制牺牲吞吐换取数据安全,异步复制反之。RocketMQ 的主从同步是"复制"层面,与刷盘(SYNC_FLUSH/ASYNC_FLUSH)是两个独立维度,需组合理解。主从切换时,若采用异步复制,丢失的是主节点尚未同步给从的消息,可通过业务重投或对账弥补。

# broker 配置主从复制模式
# brokerRole=SYNC_MASTER 或 ASYNC_MASTER
# 查看 broker 状态
sh mqadmin clusterList -n localhost:9876
#
★★★

2. RocketMQ 架构组件中 NameServer、Broker(Master/Slave)、Producer/Consumer 的职责与部署形态

请说明 RocketMQ 的 NameServer、Broker(Master/Slave)、Producer/Consumer 各组件职责与部署形态?

  • NameServer 的职责与无状态设计
  • Broker 主从与存储
  • Producer/Consumer 的职责

RocketMQ 由 NameServer、Broker、Producer、Consumer 四类角色组成。NameServer 是无状态的路由中心,负责管理 Broker 的注册信息(topic 与 broker 地址映射),Producer/Consumer 通过它获取路由信息,NameServer 之间互相独立、不共享数据,通过 Broker 心跳保持信息。Broker 是存储与服务的核心,分为 Master 与 Slave(主从部署),负责消息的存储(commitlog)、主从复制、消费进度管理;Broker 向所有 NameServer 注册并定期心跳。Producer 负责生产消息,通过 NameServer 获取路由后写入 Broker;Consumer 负责消费消息并维护消费进度。部署形态上,NameServer 通常部署 2-4 个(无状态、可多活),Broker 按主从或 Dledger 多副本部署,一个 Broker 可承载多个 topic 的分区(队列)。

NameServer 无状态是 RocketMQ 高可用的关键,它不参与消息存储,故障只影响路由发现,可通过负载均衡恢复。Broker 主从提供读写容灾,Producer/Consumer 通过 NameServer 路由寻址。整体架构比 Kafka 简单,去中心化程度较高。

# 启动 NameServer
nohup sh mqnamesrv > /data/namesrv.log 2>&1 &
# 启动 Broker(带主从配置)
nohup sh mqbroker -c broker-master.conf > /data/broker.log 2>&1 &
#
★★★

3. 消息可靠性中发送重试、刷盘策略(SYNC_FLUSH/ASYNC_FLUSH)、消费重试与死信队列(DLQ)机制

RocketMQ 如何通过发送重试、刷盘策略、消费重试与死信队列来保证消息可靠性?

  • 发送重试机制
  • 刷盘策略(SYNC_FLUSH/ASYNC_FLUSH)
  • 死信队列(DLQ)机制

发送可靠性:生产者发送失败时自动重试(默认重试次数),并可通过设置 Storage 与刷盘策略保证消息落盘。刷盘策略:SYNC_FLUSH 在消息写入磁盘后才返回成功,可靠性高但性能低;ASYNC_FLUSH 异步刷盘,性能高但 Broker 宕机时可能丢未刷盘数据。消费可靠性:Consumer 消费失败后,消息进入重试队列(retry topic,按重试次数递增延迟),多次重试仍失败则进入死信队列(DLQ,%DLQ%groupName)。运维上需监控重试次数与死信队列,合理治理死信消息(人工重投或补数)。

RocketMQ 可靠性分层:发送端(重试)、存储端(刷盘)、消费端(重试 + DLQ)。SYNC_FLUSH 与 SYNC_MASTER 组合提供最高可靠性,但吞吐受限。消费重试与 DLQ 是"至少一次"语义下的兜底:消费失败不丢消息,重试耗尽后进入 DLQ 供人工处理。运维需关注 DLQ 堆积与重试风暴。

# broker 配置刷盘策略
# flushDiskType=SYNC_FLUSH 或 ASYNC_FLUSH
# 查看 topic 与死信队列
sh mqadmin topicList -n localhost:9876
#
★★

4. NameServer 运维中无状态设计、Broker 注册与路由发现、故障影响与脑裂防护

NameServer 的无状态设计如何支撑 Broker 注册与路由发现,其故障影响与脑裂防护如何考虑?

  • NameServer 无状态设计
  • Broker 注册与心跳
  • 故障影响与脑裂防护

NameServer 是无状态的服务,所有 NameServer 实例之间不共享数据、不互相通信,每个实例独立维护一份路由表(topic 到 broker 的映射)。Broker 启动时向所有 NameServer 注册,并周期发送心跳(每 30 秒),NameServer 通过心跳超时(默认 2 分钟)判定 Broker 下线并移除路由。Producer/Consumer 启动时从 NameServer 拉取路由并缓存,也会定期刷新。NameServer 故障不影响已建立连接的消息收发(仅影响新路由发现),因为无状态、可多实例部署,通常部署多个实现容错。脑裂防护上,NameServer 无状态使其天然无脑裂问题;Broker 侧则需通过主从与 Dledger 的 Raft 机制避免多主。

NameServer 无状态是 RocketMQ 高可用的核心,故障只影响路由发现、不影响消息传输,且多实例可横向扩展。脑裂风险主要在 Broker 主从侧,通过 brokerId 与 Dledger 选举解决。运维上应监控 NameServer 与 Broker 的心跳、路由一致性。

# 查看 NameServer 中的路由信息
sh mqadmin clusterList -n localhost:9876
# 查看 broker 状态
sh mqadmin brokerStatus -n localhost:9876 -b localhost:10911
#
★★

5. 事务消息中半消息、事务回查与本地事务表机制以及未决事务如何监控与处理

RocketMQ 事务消息如何通过半消息、事务回查与本地事务表实现,未决事务如何监控与处理?

  • 半消息(half message)机制
  • 事务回查(transaction check)
  • 本地事务表

RocketMQ 事务消息用于解决"本地事务与消息发送的一致性"。流程:生产者先发送半消息(half message,对消费者不可见),执行本地事务;若本地事务提交则向 Broker 发送 commit,否则 rollback(半消息被删除)。若本地事务结果迟迟未上报,Broker 会发起事务回查(check),询问生产者本地事务的最终结果,生产者依据本地事务表(记录事务状态)返回 commit/rollback。这样保证"本地事务成功则消息必达,本地事务回滚则消息不投递"。运维上需监控未决(半消息)数量,定位事务回查超时或本地事务表不一致的问题,并处理滞留的半消息。

事务消息的本质是"两阶段 + 回查兜底"。本地事务表是回查的依据,若本地事务已提交但半消息未 commit,回查会补上 commit,保证一致性。运维关注半消息堆积,说明有事务长时间未决,需排查生产者状态或本地事务表。

# 查看未决事务消息(半消息)
# 通过 mqadmin 或控制台查看半消息队列
sh mqadmin queryMsgByKey -n localhost:9876 -t test_tx -k <key>
#
★★

6. 存储结构中 commitlog/consumequeue/index 文件、文件保留策略与容量规划

RocketMQ 的 commitlog、consumequeue、index 文件如何组织存储,文件保留策略与容量规划如何设计?

  • commitlog 顺序写存储
  • consumequeue 消费队列
  • index 索引文件

RocketMQ 的存储以 commitlog 为核心:所有消息顺序写入一个全局的 commitlog 文件(顺序写,提升性能),每个 topic 的消费队列由 consumequeue 文件组成,consumequeue 记录消息在 commitlog 中的偏移量(逻辑索引),消费时通过 consumequeue 定位到 commitlog。index 文件用于按消息 key 或时间快速检索消息。文件保留策略:通过 fileReservedTime(默认 72 小时)、deleteWhen 与磁盘空间(diskMaxUsedSpaceRatio)控制删除,超时的 commitlog 与对应的 consumequeue 会被清理。容量规划需考虑消息量、保留时长、磁盘空间与写入峰值,预留缓冲。

commitlog 顺序写是 RocketMQ 高性能的核心,consumequeue 是逻辑索引,二者分离使得写入与消费解耦。保留策略按时间与磁盘水线控制,磁盘使用率过高时 RocketMQ 会拒绝写入或触发清理。容量规划要留足余量并配合监控。

# broker 配置保留策略
# fileReservedTime=72, deleteWhen=04, diskMaxUsedSpaceRatio=75
# 查看存储目录
ls -lh /data/rocketmq/store/commitlog/ | head
#
★★

7. 常见故障排障中刷盘慢、page cache 压力、Broker 假死、路由不一致的定位

面对刷盘慢、page cache 压力、Broker 假死、路由不一致等常见故障,如何定位与处置?

  • 刷盘慢与磁盘 IO 瓶颈
  • page cache 压力与内存
  • Broker 假死的判定

刷盘慢:多由磁盘 IO 瓶颈(HDD、IO 竞争、磁盘满)导致,可通过 iostat 查看磁盘 IO 等待与吞吐,优化磁盘或减小刷盘频率。page cache 压力:RocketMQ 重度依赖 page cache,内存不足或缓存被回收会导致读写变慢,需监控内存使用与 page cache 命中率,合理配置 JVM 堆内存以留出足够 page cache。Broker 假死:Broker 进程存活但无响应(如长时间 GC、线程阻塞、磁盘卡死),会导致客户端超时、心跳中断,需通过日志、线程 dump、GC 日志定位,必要时重启或调整参数。路由不一致:NameServer 与 Broker 的路由信息不同步,导致客户端找不到队列或写入失败,需检查心跳与注册、NameServer 间的路由一致性。

排障要"看现象定位根因"。刷盘慢、page cache 压力、Broker 假死都是资源或运行时问题,需结合系统指标(IO、内存、GC)与 RocketMQ 日志日志综合判断。路由不一致多与心跳丢失或 NameServer 故障有关,重启或刷新路由可恢复。

# 查看磁盘 IO 与内存
iostat -x 1
free -m
# 查看 broker 日志
tail -n 100 /data/rocketmq/logs/broker.log
#
★★

8. 延迟/定时消息中延迟级别、时间轮实现与运维限制(延迟上限、存储成本)

RocketMQ 的延迟/定时消息如何通过延迟级别与时间轮实现,运维上有哪些限制(延迟上限、存储成本)?

  • 延迟级别(delay level)
  • 时间轮实现
  • 延迟上限与存储成本

RocketMQ 延迟消息通过延迟级别实现:消息发送时指定延迟级别(delayTimeLevel),对应固定的延迟时间(如 1s、5s、10s、30s、1min…… 18 个级别),Broker 将延迟消息放入对应级别的时间轮(timer wheel)中,时间到后再投递到真实队列。RocketMQ 5.x 支持任意延迟时间(通过定时消息的时间轮)。运维限制:延迟级别是固定的(无法任意设置秒数),延迟时间有上限(如 2 小时级别);延迟消息占用存储(在半消息/定时消息队列中),增多会增大存储成本与复杂度;需理性规划延迟场景,避免滥用。

延迟级别是预定义的固定档位,时间轮以 O(1) 复杂度调度到期消息。定时消息的存储成本在于消息要在 Broker 中暂存到延迟时间到达,期间占用磁盘。运维上需评估延迟场景的规模与存储占用,避免积压。

# 发送延迟消息(Java 客户端)
# message.setDelayTimeLevel(3);  // 对应 10s
# 查看定时消息相关队列
sh mqadmin topicList -n localhost:9876
#
★★

9. 消息堆积治理中消费 lag 监控、消费端扩容、队列再平衡与消费限流

RocketMQ 如何监控消费 lag、进行消费端扩容、队列再平衡与消费限流来治理消息堆积?

  • 消费 lag 监控
  • 消费端扩容
  • 队列再平衡

消费 lag 是消息消费进度与最新消息的差值,可通过 mqadmin consumerProgress 或控制台监控,反映堆积程度。堆积治理:先定位是消费端产能不足还是下游慢,再扩容消费端(增加消费者实例,集群模式下 RocketMQ 会自动在同一消费组内重新分配队列,实现再平衡);若队列数不足,可增加 topic 队列数(扩容队列)提升并行度。消费限流:通过消费端速率控制或下游限流,避免消费过快压垮下游。运维上需建立 lag 告警,区分"短时积压"与"持续消费不动",并评估扩容效果。

堆积的本质是产能不足。RocketMQ 集群消费模式下,同一消费组内消费者自动分配队列,增加消费者实例即可再平衡。队列数扩容需注意:扩容后重新分配,且影响消费顺序。消费限流是为了保护下游,需与扩容配合。

# 查看消费进度与 lag
sh mqadmin consumerProgress -n localhost:9876 -g my_group
# 查看消费者实例
sh mqadmin consumerStatus -n localhost:9876 -g my_group
#
★★

10. 监控与告警中 Broker 吞吐、堆积量、发送失败率、主从同步延迟的指标建设

RocketMQ 监控应建设哪些指标,Broker 吞吐、堆积量、发送失败率、主从同步延迟如何设置告警?

  • Broker 吞吐指标
  • 堆积量(lag)指标
  • 发送失败率指标

核心监控指标:Broker 吞吐(生产/消费 TPS、字节速率)、消息堆积量(各消费组 lag)、发送失败率(生产失败次数/比例)、主从同步延迟(Master 与 Slave 的 offset 差)、Broker 存活与网络、磁盘使用率。这些指标可通过 mqadmin 命令、RocketMQ 的 JMX 或 Prometheus exporter 采集。告警分级:Broker 宕机、发送失败率突增、堆积量持续增长、主从延迟过大为高优先;吞吐波动、磁盘接近水线为中优先。监控还应关注消费组 lag 的绝对值与增速,避免告警风暴。

RocketMQ 监控看"能否接收、能否消费、数据是否安全"。发送失败率与主从延迟是数据安全指标,堆积量是业务健康指标,吞吐是性能指标。告警需设置合理阈值与持续时间,配合消费 lag 与主从延迟防止误报。

# 查看 broker 运行时指标
sh mqadmin brokerRuntime -n localhost:9876 -b localhost:10911
# 查看消费 lag
sh mqadmin consumerProgress -n localhost:9876 -g my_group
#
★★

11. 集群扩缩容中 Broker 增减、Topic 队列扩缩、流量迁移与优雅下线

RocketMQ 集群如何实施 Broker 增减、Topic 队列扩缩、流量迁移与优雅下线?

  • Broker 增加与减少
  • Topic 队列数量调整
  • 流量迁移

Broker 增加:启动新 Broker 并注册到 NameServer,新流量可路由到新 Broker;可通过迁移 topic 或调整 topic 的队列分布来均衡流量。Broker 减少(下线):下线前需确认无消费者依赖,先迁移该 Broker 上的 topic 队列或确认主从 OK,再优雅停服,避免消息丢失。Topic 队列扩缩:通过 updateTopic 调整队列数(写队列与读队列),扩容提升并行度,缩容需谨慎(可能影响消费)。流量迁移:通过调整 topic 队列分配、控制客户端路由或使用工具将流量从旧 Broker 迁移到新 Broker。优雅下线:先停止接收新消息(如下线写队列),等消费完积压后再停服。

扩缩容要"平滑、可回退"。Broker 下线是高风险操作,需保证主从切换与数据不丢。队列扩缩影响消费分配与顺序,需在业务低峰执行。流量迁移要渐进,观察对账后再继续。

# 调整 topic 队列数
sh mqadmin updateTopic -n localhost:9876 -b localhost:10911 \
  -t orders -r 8 -w 8
# 关闭 broker
sh mqshutdown broker
#
★★

12. 顺序消息中全局/分区顺序(MessageQueueSelector)的实现与顺序消费故障处理

RocketMQ 如何通过 MessageQueueSelector 实现全局/分区顺序消息,顺序消费故障如何处理?

  • 顺序消息的两种级别(全局/分区)
  • MessageQueueSelector 实现
  • 顺序消费的机制

RocketMQ 顺序消息分为全局顺序(所有消息按序)与分区顺序(同一 key 的消息按序)。实现依赖 MessageQueueSelector:生产者根据业务 key 通过选择器将同一 key 的消息路由到同一个队列,从而保证该队列内消息有序。消费端使用 MessageQueueListener(顺序消费)单线程按序消费该队列,且消费失败会用重试机制但不会跳过。顺序消费故障处理:一旦某条消息消费失败,顺序消费会阻塞该队列后续消息,可能导致该队列整体暂停,需快速定位失败原因并修复,避免长时间阻塞;监控队列消费进度与失败率。

顺序保证的本质是"同一队列内有序 + 单线程消费"。MessageQueueSelector 保证同一 key 落在同一队列,顺序消费保证队列内逐条处理。故障时顺序消费为避免乱序会阻塞后续消息,因此需快速排障,否则整队列停滞。

// 生产者按业务 key 选择队列
SendResult result = producer.send(msg, new MessageQueueSelector() {
    public MessageQueue select(List<MessageQueue> mqs, Message msg, Object arg) {
        int id = String.valueOf(arg).hashCode() % mqs.size();
        return mqs.get(id);
    }
}, orderId);
#
★★

13. 集群消费与广播消费中两种消费模式的语义差异、适用场景与运维上的坑(广播模式下的堆积与重置)?

RocketMQ 集群消费与广播消费在语义、适用场景上有何差异,广播模式下有哪些运维坑(堆积与重置)?

  • 集群消费与广播消费的语义
  • 适用场景
  • 广播模式下的堆积

集群消费(CLUSTERING):同一消费组内,一条消息只被一个消费者实例消费,消息被分摊到组内各实例,适合工作队列、负载均衡场景。广播消费(BROADCASTING):同一消费组内,每个消费者实例都会消费同一条消息,适合所有实例都需要的场景(如配置下发、缓存更新)。运维坑:广播模式下每个实例各自维护消费进度,若某实例新上线或消费进度丢失,会重新从最新或初始位置消费,可能造成消息重复或堆积;广播模式下消费进度按实例独立,无法像集群模式那样通过组内再平衡分摊,扩容需谨慎;广播消费的进度重置(从最后消费还是从头)需明确,避免误读历史数据。

集群消费是"一条消息被消费一次",广播消费是"每条消息每个实例都消费"。广播模式没有实例间分摊,每个实例独立消费同一份数据,因此实例越多,重复消费的任务越多,且消费进度独立管理,容易产生进度不一致或重置导致的重读。运维需明确广播模式的进度语义。

# 消费者设置消费模式(Java)
# consumer.setMessageModel(MessageModel.BROADCASTING);
# 广播模式:每个实例独立存进度;集群模式:统一存 Broker
#
★★

14. 消息过滤中 Tag 过滤与 SQL92 属性过滤的实现位置与性能差异以及如何使用属性过滤避免 Tag 爆炸?

RocketMQ 的 Tag 过滤与 SQL92 属性过滤在实现位置与性能上有何差异,如何用属性过滤避免 Tag 爆炸?

  • Tag 过滤的实现位置
  • SQL92 属性过滤的实现位置
  • 两者的性能差异

Tag 过滤在 Broker 端进行:消息带 Tag,消费者订阅时指定 Tag,Broker 在分发时按 Tag 过滤,只投递匹配的 Tag,性能高、开销小。SQL92 属性过滤:消费者通过 SQL 表达式(如 a > 100 AND b = 'x')基于消息的扩展属性过滤,Broker 在消费端投递前计算表达式,支持更复杂的过滤逻辑,但性能开销高于 Tag 过滤(需解析表达式)。当 Tag 数量爆炸(业务 Tag 过多)时,Tag 过滤难以维护且影响性能,可改用 SQL92 属性过滤,把过滤条件放到属性字段上,减少 Tag 数量与管理成本。

Tag 过滤是"提前在 Broker 分发时过滤",简单高效;SQL92 过滤是"按属性表达式计算",灵活但开销大。Tag 爆炸是 Tag 数量过多导致管理复杂、订阅匹配困难,用属性过滤把细分维度放到属性字段,让 Tag 保持少量,是治理手段。

// SQL92 属性过滤
consumer.subscribe("orders", MessageSelector.bySql("price > 100 AND region = 'cn'"));
#
★★

15. 消费幂等中消息重投与重复消费的必然性以及去重表/业务状态机等幂等方案的实现要点?

RocketMQ 消息重投与重复消费为何必然发生,去重表/业务状态机等幂等方案如何实现?

  • 重复消费的必然性
  • 去重表方案
  • 业务状态机方案

RocketMQ 采用 at-least-once 语义,消费端可能因重试、主从切换、客户端重启等导致同一消息被重复消费,因此消费端必须实现幂等。幂等方案:去重表(用消息唯一键如业务 ID 建唯一索引,消费前先查重,命中则跳过,否则插入并处理,配合事务保证原子性);业务状态机(以业务状态为幂等依据,若状态已达成目标则跳过,如订单状态流转到已支付则不再重复扣款);幂等键 + 分布式锁。实现要点:唯一键要稳定、去重与业务处理要原子(同事务或先落库再处理)、处理失败要能重试而不重复生效。

重复消费在 at-least-once 下不可避免,幂等是对"消费副作用"的防护。去重表强调"键唯一预防",状态机强调"业务状态判断",两者都要求原子性。实现要点是使"重复执行与执行一次"效果一致,这是幂等定义。

// 去重表幂等(伪代码)
if (!dedupTable.exists(orderId)) {
    // 插入去重记录 + 业务处理,同一事务
    dedupTable.insert(orderId);
    process(orderId);
}
#
★★

16. 4.x 到 5.x 迁移中 Broker 职责拆分、gRPC 协议与 Proxy 层引入后的部署差异与兼容性?

RocketMQ 从 4.x 迁移到 5.x,Broker 职责拆分、gRPC 协议与 Proxy 层引入后部署差异与兼容性如何?

  • 4.x 到 5.x 的架构变化
  • Broker 职责拆分
  • gRPC 协议与 Proxy 层

RocketMQ 5.x 引入了新的架构:统一的代理层(Proxy)与 gRPC 协议,客户端可通过 gRPC 协议访问,同时兼容旧版 Remoting 协议。5.x 中 Broker 职责有所拆分,引入 controller(基于 Raft)管理元数据与主从切换,Broker 的职责更聚焦于存储与消息服务。部署差异:5.x 支持 Proxy 模式(客户端连 Proxy,Proxy 转 Broker)与无 Proxy 模式(客户端直连 Broker,兼容旧版),Proxy 提供协议转换、负载均衡与治理能力。兼容性:5.x 客户端与 4.x Broker 兼容(通过协议适配),但 5.x 的 gRPC 特性需 5.x Broker 支持;迁移需灰度验证客户端与 Broker 的版本兼容矩阵。

5.x 的核心变化是"代理层 + gRPC + controller 高可用"。Proxy 层让集群具备更便捷的治理与协议扩展能力,但引入额外一跳。Broker 职责拆分与 controller 提升高可用自动化。迁移要关注版本兼容与协议差异,避免功能失效。

# 5.x 启动 proxy(可选项)
sh mqproxy -n localhost:9876
# 客户端通过 gRPC 或 Remoting 连接
#

17. 与 Kafka 的对比与迁移中消费/存储模型差异、迁移工具与双跑验证

RocketMQ 与 Kafka 在消费/存储模型上有何差异,RocketMQ 迁移到 Kafka(或反之)如何通过工具与双跑验证?

  • 消费/存储模型差异
  • 迁移工具
  • 双跑验证

存储模型:Kafka 以分区(partition)为单位,每个分区独立日志,顺序写;RocketMQ 以全局 commitlog 存储,consumequeue 作为逻辑索引,二者写入路径不同。消费模型:Kafka 消费组按分区分配,顺序性在分区内;RocketMQ 集群消费按队列分配,顺序性在队列内。两者都支持消息堆积与消费组。迁移工具:可借助消息桥接(如 MirrorMaker for Kafka、RocketMQ 的跨 MQ 桥接工具)或自研双写;迁移策略:双写(同时写两边)→ 双跑(两边都消费,对比对账)→ 切换消费 → 停机。双跑验证要对比消息数、消息内容、消费 lag 与业务结果,确保一致后再切换。

两者模型差异导致"迁移到对方"需处理分区/队列对齐、顺序语义、消费进度映射。迁移高风险,双写 + 双跑 + 对账是标准路径,保留回退。若业务强依赖分区顺序,需评估 RocketMQ 队列与 Kafka 分区的映射。

# 双写:应用同时写两个 MQ
# 对账:对比两边消息数与内容
# 迁移工具:自研桥接 或 第三方 MQ 桥接
echo "迁移路径:双写 -> 双跑对账 -> 切换消费 -> 停机"
#

18. 多副本与容灾中 Dledger(Raft)模式、跨机房复制与主从切换

RocketMQ 如何通过 Dledger(Raft)模式实现多副本容灾,跨机房复制与主从切换如何设计?

  • Dledger(Raft)模式
  • 跨机房复制
  • 主从切换

Dledger 模式让 RocketMQ 的 Broker 通过 Raft 协议实现多副本(如 3 副本),副本间自动选举 leader,leader 写入成功后多数副本确认才返回,实现数据安全与自动故障切换。相比传统主从(Master/Slave)模式,Dledger 无需人工指定主从,leader 故障时自动从 follower 中选出新 leader,避免"切换主从"的手工操作。跨机房复制:通过 RocketMQ 的跨机房同步(如多活复制、正则复制工具)或基于 Dledger 的跨机房多副本,实现容灾。主从切换:Dledger 自动选举,传统主从可手动切换(brokerId 0 为主)。运维需评估 Raft 的多数派副本与跨机房延迟。

Dledger 是 RocketMQ 高可用自动化的重要组件,用 Raft 解决"多副本一致性 + 自动选主"。跨机房复制要权衡延迟与一致性,多数派副本需在跨机房满足。运维上监控副本状态、leader 健康与同步延迟。

# Dledger 模式 broker 配置
# enableDLegerCommitLog=true
# dLegerGroup=group1, dLegerPeers=n0-127.0.0.1:40911;n1-...
# 查看 broker 健康
sh mqadmin clusterList -n localhost:9876
#

19. 安全与运维工具中 ACL、消息轨迹(trace)、消息查询与回放的使用

RocketMQ 的 ACL、消息轨迹(trace)、消息查询与回放等安全与运维工具如何部署使用?

  • ACL 权限控制
  • 消息轨迹(trace)
  • 消息查询

ACL:RocketMQ 通过配置文件(plain_acl.yml)定义用户与权限,支持 Topic/Group 粒度的读写权限,控制客户端访问。消息轨迹(trace):开启 trace 后,RocketMQ 收集消息的发送、存储、消费各环节时间戳,用于排查消息去向与延迟,常通过 trace topic 存储。消息查询:通过 mqadmin queryMsgByKey、queryMsgById 按 key 或 id 查询消息,定位消息内容与状态。消息回放:基于查询到的消息重新投递,用于补偿与对账。运维上需合理配置 ACL 权限、开启 trace 并监控 trace 数据量,避免 trace 本身造成堆积。

安全与可观测是运维的左右手。ACL 保障权限边界,trace 提供全链路可观测,查询与回放用于故障定位与补偿。三者的组合让"消息去哪了、谁在访问、如何恢复"都有依据。运维需权衡 trace 的存储成本与可观测价值。

# 按 key 查询消息
sh mqadmin queryMsgByKey -n localhost:9876 -t orders -k <key>
# 按 id 查询消息
sh mqadmin queryMsgById -n localhost:9876 -i <msgId>
#

20. 性能与 JVM 调优中堆外内存、page cache 与刷盘策略的关系以及 Broker JVM 参数设置的要点?

RocketMQ Broker 的性能与 JVM 调优中,堆外内存、page cache 与刷盘策略的关系如何,JVM 参数设置要点是什么?

  • 堆外内存与 page cache 的关系
  • 刷盘策略与内存
  • Broker JVM 参数要点

RocketMQ 的读写依赖操作系统 page cache:消息写入先进入 page cache(堆外内存),由内核异步刷盘;读也从 page cache 命中。因此 JVM 堆内存不宜过大,要留出足够内存给 page cache,否则缓存压力大导致读写变慢。刷盘策略(SYNC_FLUSH/ASYNC_FLUSH)决定多快把 page cache 刷到磁盘,SYNC_FLUSH 更安全但更慢。JVM 参数要点:堆内存(Xmx)适中(如 8-16G,视机器内存),避免堆过大挤占 page cache;合理设置 GC(如 G1)、避免频繁 Full GC;控制线程池与连接数。调优目标是"堆与 page cache 平衡 + 刷盘与吞吐平衡"。

RocketMQ 高性能依赖"堆外 page cache + 顺序写"。JVM 堆过大反而挤占 page cache,是常见误区。刷盘策略是安全与性能的权衡。调优需结合机器内存、消息量、可靠性要求综合设置。

# 典型 JVM 参数(不挤占 page cache)
# -Xms8g -Xmx8g -XX:+UseG1GC
# 手动调整 broker 的 JVM 内存(编辑 bin/runbroker.sh 的 JAVA_OPT)