基本最佳实践
生产者 (Producer)
发送消息的注意事项
Tags 的使用
应用可以被标识为 Topic,消息子类型可以被标识为 Tags。Tags 可以由应用自由设置。只有当生产者在发送消息时设置了 Tags,消费者在订阅消息时才能通过 Broker 使用 Tags 过滤消息。5.x SDK 可调用 messageBuilder.setTag("messageTag"),历史版本可调用 message.setTags("messageTag")。
Keys 的使用
在业务层面,建议为每条消息映射一个唯一的业务标识并设置到 keys 字段中,以便在未来定位消息丢失问题。服务端会为每条消息创建索引(哈希索引),应用可以通过 Topic 和 Key 查询消息内容,以及该消息被谁消费。由于是哈希索引,请确保 Key 尽可能唯一,以避免潜在的哈希碰撞。常见的设置策略是使用订单 ID、用户 ID、请求 ID 等离散的唯一标识符。
打印日志
无论消息发送成功还是失败,都需要打印消息日志以便进行服务故障排查。只要没有抛出异常,就表示消息发送成功。
消息发送失败的处理方法
生产者本身的 send 方法支持内部重试,5.x 重试逻辑参考 发送重试策略:
上述策略在一定程度上保证了消息发送的成功率。如果业务要求消息发送不丢失,仍需覆盖可能出现的异常情况。例如,在调用同步发送方法且发送失败时,尝试将消息存储到数据库中,并由后台线程定期重试,以确保消息到达 Broker。
上述数据库重试方法未集成到 MQ 客户端中,而是要求应用自行完成,主要是基于以下考虑:第一,MQ 客户端被设计为无状态模式,方便任意水平扩展,机器资源消耗仅限于 CPU、内存、网络。第二,如果 MQ 客户端内部集成了 KV 存储模块,只有同步落盘才能保证数据可靠,而同步落盘本身有较大的性能开销,因此通常使用异步落盘。且由于应用关闭过程不受 MQ 运维人员控制,经常会出现 kill -9 这种强制关闭,导致数据未能及时落盘而丢失。第三,生产者所在的机器可靠性较低,通常是虚拟机,不适合存储重要数据。综上所述,建议重试过程由应用控制。
消费者 (Consumer)
消费过程幂等性
RocketMQ 无法避免消息重复(Exactly Once),因此如果业务对消费重复非常敏感,在业务层面进行去重至关重要。这可以借助关系型数据库来实现。首先需要确定消息的唯一键,可以是 msgId 或消息内容中的唯一标识字段(如订单 ID)。在消费前,先确定唯一键在关系型数据库中是否存在。如果不存在,则插入并消费;否则,直接跳过。(实际处理过程应考虑原子性问题,判断是否存在主键冲突,如果插入失败,直接跳过)
MsgId 必须是全局唯一标识符,但在实际操作中,可能会出现同一消息有两个不同 msgId 的情况(消费者主动重传、客户端投递机制导致的重复等),因此必须基于业务字段进行去重。
慢消费处理
提高消费并行度
绝大多数消息消费是 IO 密集型的,即可能在操作数据库或调用 RPC,这类消费的速率取决于后端数据库或外部系统的吞吐量。通过提高消费并行度可以提高总消费吞吐量,但当并行度增加到一定程度时,性能会下降。因此,应用必须设置合理的并行度。修改消费并行度有几种方式:
- 在同一个 ConsumerGroup 中,增加消费者实例的数量以提高并行度(注意:超过订阅队列数的消费者实例无效)。您可以增加机器,或在现有机器上启动多个进程。
- 提高单个消费者的消费并行线程数。5.x PushConsumer SDK 可以通过 PushConsumerBuilder.setConsumptionThreadCount() 设置线程数;SimpleConsumer 可以自由地从业务线程增加并发,底层是线程安全的;历史 SDK PushConsumer 可以通过修改参数 consumeThreadMin 和 consumeThreadMax 来实现。
批量消费
如果某些业务流程支持批量消费,消费吞吐量可以得到极大提升。例如,订单扣减应用处理一个订单需要 1 秒,而一次处理 10 个订单可能只需要 2 秒,这样消费吞吐量就得到了大幅提升。建议使用 5.x SDK 的 SimpleConsumer,设置每次接口调用的批量大小,一次拉取多条消息。
重置位点以跳过不重要的消息
在消息堆积的情况下,如果消费速率跟不上投递速率,且业务对数据完整性要求不高,可以选择丢弃不重要的消息。建议使用重置位点功能,直接将消费位点调整到指定的时间或位置。
优化单条消息消费过程
例如,一条消息的消费过程如下:
- 查询[数据 1]根据消息从数据库查询
- 查询[数据 2]根据消息从数据库查询
- 复杂的业务计算
- 插入[数据 3]到数据库
- 插入[数据 4]到数据库
在此消息消费过程中有四次与数据库的交互。如果计算每次交互耗时 5ms,总时间为 20ms。假设业务计算耗时 5ms,总耗时为 25ms。因此,如果能将四次数据库交互优化为两次,总时间可以优化到 15ms,这意味着整体性能提高了 40%。因此,如果应用对延迟敏感,可以将数据库部署在 SSD 磁盘上。与 SCSI 磁盘相比,前者的响应时间 (RT) 要小得多。
消费打印日志
如果消息数量较少,建议在消费入口方法中打印消息,以便记录耗时较长的消费。
new MessageListener() {
@Override
public ConsumeResult consume(MessageView messageView) {
LOGGER.info("Consume message={}", messageView);
//Do your consume process
return ConsumeResult.SUCCESS;
}
}
如果能够打印每条消息的消费耗时,对于排查线上慢消费等问题会更加方便。
Broker
Broker 角色
Broker 角色分为 ASYNC_MASTER、SYNC_MASTER 和 SLAVE。如果对消息可靠性有严格要求,请部署 SYNC_MASTER 加 SLAVE。如果不需要消息可靠性,则部署 ASYNC_MASTER 加 SLAVE。如果仅为方便测试,可以选择仅部署 ASYNC_MASTER 或 SYNC_MASTER。
刷盘方式 (FlushDiskType)
与 ASYNC_FLUSH 相比,SYNC_FLUSH 会损失性能,但更可靠。因此,必须根据实际业务场景进行权衡。
Broker 配置
| 参数 | 默认值 | 描述 |
|---|---|---|
| listenPort | 10911 | 接受客户端连接的监听端口 |
| namesrvAddr | null | NameServer 地址 |
| brokerIP1 | 当前网络地址 | Broker 当前监听的 IP 地址 |
| brokerIP2 | 同 brokerIP1 | 当存在主从 Broker 时,如果 Broker 主节点配置了 brokerIP2 属性,从节点将连接到主节点配置的 brokerIP2 进行同步 |
| brokerName | null | Broker 名称 |
| brokerClusterName | DefaultCluster | 该 Broker 所属的集群名称 |
| brokerId | 0 | Broker ID,0 表示 Master,其他正整数表示 Slave |
| storePathCommitLog | $HOME/store/commitlog/ | 存储 Commit Log 的路径 |
| storePathConsumerQueue | $HOME/store/consumequeue/ | 存储消费队列的路径 |
| mappedFileSizeCommitLog | 1024 * 1024 *1024(1G) | Commit Log 映射文件大小 |
| deleteWhen | 04 | 在一天中的什么时间删除超过保留时间的文件 |
| fileReservedTime | 72 | 文件保留时间(小时) |
| brokerRole | ASYNC_MASTER | SYNC_MASTER/ASYNC_MASTER/SLAVE |
| flushDiskType | ASYNC_FLUSH | SYNC_FLUSH/ASYNC_FLUSH。SYNC_FLUSH 模式的 Broker 保证在接收到生产者确认之前刷盘。ASYNC_FLUSH Broker 使用异步刷盘方式以获得更好的性能。 |
Broker 日志管理
Broker 的默认日志路径位于 ${user.home}/logs/rocketmqlogs/。您可以通过编辑二进制包 conf 目录下的 xx.logback.xml 文件来修改日志级别和路径。
注意:请确保您的日志已妥善保护,以防止敏感信息泄露。