这是《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
生产检查清单
- Topic、队列数、权限已评审;
- 副本与刷盘策略满足 RPO;
- 磁盘保留覆盖最大 lag 和恢复窗口;
- ACL 账号最小权限;
- 生产端有 Outbox 或事务补偿;
- 消费端幂等;
- 重试与死信治理入口可用;
- 消息轨迹和业务键可查询;
- 监控告警覆盖集群、Topic、消费组;
- 扩缩容和切主已演练;
- 备份可恢复;
- 迁移有对账和回滚方案。
版本注意事项
| 能力 | 注意 |
|---|---|
| 定时消息 | 5.x 能力与 4.x 固定级别差异较大 |
| Proxy | 客户端协议和 SDK 版本需匹配 |
| Controller | 部署与配置随版本变化 |
| Dashboard | 功能与 API 随版本变化 |
| 事务消息 | 不同客户端构造函数和监听器 API 有差异 |
思考题
- 你的核心 Topic 使用哪种刷盘和复制组合?
- 消费组最大可接受消息年龄是多少?
- 死信重放入口在哪里?
- 位点重置是否已有审批和审计?
- 最近一次备份恢复演练是什么时候?