这是《RocketMQ 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。 普通消息是最常用的消息模型:生产者把消息发送到某个 Topic 的一个队列,消费者按消费组拉取并处理。它适合事件通知、异步解耦、削峰填谷,不承诺同一业务 key 的全局顺序。
7.1 发送流程
Producer
-> 从 NameServer 获取 Topic 路由
-> 选择 Broker 和 Queue
-> 构造消息
-> 写入 Broker
-> 返回 SendResult
完整发送示例:
DefaultMQProducer producer = new DefaultMQProducer("order-producer-group");
producer.setNamesrvAddr("rocketmq-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");
message.putUserProperty("source", "order-service");
SendResult result = producer.send(message);
System.out.printf("status=%s, msgId=%s, queue=%s%n",
result.getSendStatus(),
result.getMsgId(),
result.getMessageQueue());
} finally {
producer.shutdown();
}
setKeys 不影响路由,但会写入索引,后续可按业务键查询消息。
7.2 消息体设计
推荐使用明确版本化的 JSON 或 Avro:
{
"eventId": "01J8ZC9Q7M6Q",
"eventType": "OrderCreated",
"version": 1,
"occurredAt": "2026-08-25T10:00:00+08:00",
"data": {
"orderNo": "O202608250001",
"userId": "U10001",
"amount": 199.00
}
}
规范:
- 必须有全局唯一
eventId; - 必须有事件类型和版本;
- 时间带时区;
- 消息体只放必要字段;
- 大对象存对象存储,消息只放引用;
- 不依赖线程上下文自动生成隐式字段。
7.3 Tag 与属性
Tag 用于订阅过滤:
consumer.subscribe("OrderTopic", "OrderCreated || OrderPaid");
User Property 用于复杂 SQL 过滤:
message.putUserProperty("region", "cn-east");
message.putUserProperty("level", "vip");
consumer.subscribe("OrderTopic", MessageSelector.bySql("region = 'cn-east' AND level = 'vip'"));
选择建议:
| 方式 | 特点 |
|---|---|
| Tag | 简单、常用、性能好 |
| SQL92 | 表达能力强,需要 Broker 开启支持 |
| 业务端过滤 | 灵活但会消耗消费带宽 |
不要把核心业务逻辑建立在过度复杂的过滤表达式上,复杂规则更适合由业务服务判断。
7.4 同步、异步和单向发送
| 模式 | 是否等待结果 | 适用 |
|---|---|---|
| sync | 是 | 重要业务事件 |
| async | 回调返回 | 链路长、并发高的场景 |
| oneway | 否 | 日志、指标等可容忍丢失 |
异步发送:
producer.send(message, new SendCallback() {
@Override
public void onSuccess(SendResult result) {
metric.success();
}
@Override
public void onException(Throwable e) {
log.error("send failed, key={}", "O202608250001", e);
metric.failure();
}
});
异步发送失败必须落盘、重投或交给补偿任务处理,不能只打印日志。
7.5 发送结果
SendStatus 常见值:
| 状态 | 含义 | 处理 |
|---|---|---|
| SEND_OK | 发送成功 | 继续业务 |
| FLUSH_DISK_TIMEOUT | 刷盘超时 | 根据刷盘策略评估风险 |
| FLUSH_SLAVE_TIMEOUT | 从节点同步超时 | 关注副本状态 |
| SLAVE_NOT_AVAILABLE | 从节点不可用 | 告警并检查集群 |
发送方不能只判断是否抛异常,还应记录 SendStatus 和消息键,便于对账。
7.6 消费普通消息
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("order-consumer-group");
consumer.setNamesrvAddr("rocketmq-namesrv:9876");
consumer.subscribe("OrderTopic", "OrderCreated");
consumer.registerMessageListener((MessageListenerConcurrently) (msgs, context) -> {
for (MessageExt msg : msgs) {
handle(msg);
}
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
});
consumer.start();
处理建议:
- 解析失败记录后确认,避免无限重试;
- 临时异常返回重试;
- 幂等键使用
eventId; - 外部调用设置超时;
- 不在消费线程里做长时间阻塞;
- 消费结果必须可观测。
7.7 丢失与重复
普通消息通常提供至少一次投递语义,生产端和消费端都可能产生重复。
降低丢失风险:
业务落库
-> 发送消息
-> 根据发送结果更新事件状态
-> 定时补偿未确认事件
关键点:
- 生产端保存事件发送状态;
- 重要消息使用同步发送;
- Broker 使用合适的刷盘和复制策略;
- 消费成功后业务确认;
- 建立消息对账;
- 消费端幂等。
7.8 性能实践
| 项目 | 建议 |
|---|---|
| 生产者 | 复用实例,避免每次创建 |
| 消息体 | 控制大小,大文件走对象存储 |
| 队列数 | 与消费并行度匹配 |
| 批量 | 只在业务允许时使用 |
| 超时 | 覆盖网络和 Broker 写入时间 |
| 重试 | 与业务可重复性匹配 |
生产者实例启动成本较高,应在应用生命周期内复用,并在关闭时调用 shutdown()。
7.9 常见故障
| 现象 | 排查 |
|---|---|
| 发送超时 | NameServer、Broker 地址、网络、磁盘、队列热点 |
| Topic 不存在 | 未创建、权限不足、环境错误 |
| 消息查不到 | 未设置 Key、时间范围错、已过期 |
| 消费乱序 | 使用了并发消费,普通消息本不保序 |
| 消息重复 | 发送重试或消费重试 |
| 消息丢失 | 生产未确认、存储策略、消费跳过 |
本章小结
普通消息的核心是可靠发送、明确消息契约、合理使用 Tag 和属性、消费端幂等。它不保证全局顺序,但胜于通用性和吞吐能力。生产环境要把发送状态、消费结果、消息轨迹和对账串联起来,才能有效控制丢失与重复。
思考题
- 同步、异步、单向发送分别适合什么业务?
- 为什么消息必须包含
eventId和版本? SendStatus为什么要持久化或打点?- 普通消息出现乱序是否是缺陷?
- 如何用补偿任务降低生产端消息丢失风险?