RocketMQNotes

附录:附录:RocketMQ 速查手册

zjc 于 2026-02-01 发布

这是《RocketMQ 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。 本速查以 RocketMQ 5.x 为主线,兼顾 4.x 常用能力。不同版本命令、参数和默认值可能变化,生产操作前以当前版本文档为准。

核心概念

概念 说明
NameServer 轻量路由中心,维护 Broker 和 Topic 路由
Broker 消息存储与服务节点
Topic 消息逻辑分类
Queue Topic 内物理并行队列
ProducerGroup 生产者逻辑分组
ConsumerGroup 消费者逻辑分组
CommitLog 所有 Topic 消息的物理顺序日志
ConsumeQueue Topic 队列逻辑索引
IndexFile Key 与时间查询索引
DLedger / Controller 副本一致性与自动高可用相关模式
Proxy 5.x 接入层,可转发或协调客户端请求

常用端口

服务 常见端口
NameServer 9876
Broker 10911、10909、10912 等随配置变化
Proxy 8080、8081 等随配置变化
Dashboard 8080

端口以实际配置文件为准。

Docker 快速启动

docker run -d --name rocketmq-namesrv apache/rocketmq:5.3.1 sh mqnamesrv
docker run -d \
  --name rocketmq-broker \
  --link rocketmq-namesrv:namesrv \
  -e "NAMESRV_ADDR=namesrv:9876" \
  apache/rocketmq:5.3.1 sh mqbroker
docker run -d \
  --name rocketmq-dashboard \
  -p 8080:8080 \
  -e "JAVA_OPTS=-Drocketmq.namesrv.addr=namesrv:9876" \
  apache/rocketmq-dashboard:latest

生产环境固定版本,不建议使用 latest

mqadmin 常用命令

mqadmin clusterList -n namesrv:9876
mqadmin topicList -n namesrv:9876
mqadmin topicRoute -n namesrv:9876 -t OrderTopic
mqadmin topicStatus -n namesrv:9876 -t OrderTopic
mqadmin consumerProgress -n namesrv:9876 -g order-consumer
mqadmin consumerConnection -n namesrv:9876 -g order-consumer
mqadmin updateTopic -n namesrv:9876 -c DefaultCluster -t OrderTopic -r 6 -w 6

高危命令:

mqadmin deleteTopic
mqadmin resetOffsetByTime
mqadmin wipeWritePermission

执行前必须确认环境、Topic、消费位点和回滚方案。

生产者模板

DefaultMQProducer producer = new DefaultMQProducer("order-producer-group");
producer.setNamesrvAddr("namesrv:9876");
producer.setSendMsgTimeout(3000);
producer.setRetryTimesWhenSendFailed(2);
producer.start();

try {
    Message message = new Message(
            "OrderTopic",
            "OrderCreated",
            "O202608250001".getBytes(StandardCharsets.UTF_8));
    message.setKeys("O202608250001");
    SendResult result = producer.send(message);
    System.out.println(result.getSendStatus() + "," + result.getMsgId());
} finally {
    producer.shutdown();
}

消费者模板

DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("order-consumer");
consumer.setNamesrvAddr("namesrv:9876");
consumer.subscribe("OrderTopic", "OrderCreated || OrderPaid");
consumer.setConsumeThreadMin(4);
consumer.setConsumeThreadMax(8);
consumer.setConsumeTimeout(15);
consumer.registerMessageListener((MessageListenerConcurrently) (msgs, context) -> {
    try {
        msgs.forEach(this::handle);
        return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
    } catch (TemporaryException e) {
        return ConsumeConcurrentlyStatus.RECONSUME_LATER;
    }
});
consumer.start();

顺序消息

producer.send(message, (queues, msg, arg) -> {
    int index = Math.floorMod(String.valueOf(arg).hashCode(), queues.size());
    return queues.get(index);
}, orderNo);
consumer.registerMessageListener((MessageListenerOrderly) (msgs, context) -> {
    try {
        msgs.forEach(this::handle);
        return ConsumeOrderlyStatus.SUCCESS;
    } catch (Exception e) {
        context.setSuspendCurrentQueueTimeMillis(1000);
        return ConsumeOrderlyStatus.SUSPEND_CURRENT_QUEUE_A_MOMENT;
    }
});

延迟消息

固定级别:

message.setDelayTimeLevel(16);

5.x 定时:

message.setDeliverTimeMs(System.currentTimeMillis() + TimeUnit.MINUTES.toMillis(30));

必须确认客户端和 Broker 均支持对应能力。

事务消息要点

half message
  -> execute local transaction
  -> COMMIT_MESSAGE / ROLLBACK_MESSAGE / UNKNOW
  -> checkback by transactionId

回查必须查询持久化状态表,不能依赖内存。

存储速查

store/commitlog
store/consumequeue/Topic/queueId
store/index
store/config
store/checkpoint
文件 作用
CommitLog 消息物理正文
ConsumeQueue 队列逻辑索引
IndexFile Key 和时间查询
checkpoint 恢复检查点
config Topic、消费位点等配置

监控指标

broker_put_tps
broker_get_tps
broker_put_latency_ms
commitlog_flush_latency_ms
dispatch_lag
disk_usage_ratio
replica_sync_diff
producer_send_failure_total
consumer_lag
consumer_retry_total
dead_letter_oldest_age_seconds

排障路径

send fail
  -> client config
  -> topic route
  -> acl
  -> broker metrics
message missing
  -> business outbox
  -> send result
  -> broker message
  -> consumer trace
  -> business result
consumer lag
  -> group online
  -> queue assignment
  -> consume latency
  -> retry / dead letter
  -> downstream health

生产检查清单

  1. Topic、队列数、权限已评审;
  2. 副本与刷盘策略满足 RPO;
  3. 磁盘保留覆盖最大 lag 和恢复窗口;
  4. ACL 账号最小权限;
  5. 生产端有 Outbox 或事务补偿;
  6. 消费端幂等;
  7. 重试与死信治理入口可用;
  8. 消息轨迹和业务键可查询;
  9. 监控告警覆盖集群、Topic、消费组;
  10. 扩缩容和切主已演练;
  11. 备份可恢复;
  12. 迁移有对账和回滚方案。

版本注意事项

能力 注意
定时消息 5.x 能力与 4.x 固定级别差异较大
Proxy 客户端协议和 SDK 版本需匹配
Controller 部署与配置随版本变化
Dashboard 功能与 API 随版本变化
事务消息 不同客户端构造函数和监听器 API 有差异

思考题

  1. 你的核心 Topic 使用哪种刷盘和复制组合?
  2. 消费组最大可接受消息年龄是多少?
  3. 死信重放入口在哪里?
  4. 位点重置是否已有审批和审计?
  5. 最近一次备份恢复演练是什么时候?