跳转至主要内容
版本: 5.0

基本最佳实践

生产者 (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 配置

参数默认值描述
listenPort10911接受客户端连接的监听端口
namesrvAddrnullNameServer 地址
brokerIP1当前网络地址Broker 当前监听的 IP 地址
brokerIP2同 brokerIP1当存在主从 Broker 时,如果 Broker 主节点配置了 brokerIP2 属性,从节点将连接到主节点配置的 brokerIP2 进行同步
brokerNamenullBroker 名称
brokerClusterNameDefaultCluster该 Broker 所属的集群名称
brokerId0Broker ID,0 表示 Master,其他正整数表示 Slave
storePathCommitLog$HOME/store/commitlog/存储 Commit Log 的路径
storePathConsumerQueue$HOME/store/consumequeue/存储消费队列的路径
mappedFileSizeCommitLog1024 * 1024 *1024(1G)Commit Log 映射文件大小
deleteWhen04在一天中的什么时间删除超过保留时间的文件
fileReservedTime72文件保留时间(小时)
brokerRoleASYNC_MASTERSYNC_MASTER/ASYNC_MASTER/SLAVE
flushDiskTypeASYNC_FLUSHSYNC_FLUSH/ASYNC_FLUSH。SYNC_FLUSH 模式的 Broker 保证在接收到生产者确认之前刷盘。ASYNC_FLUSH Broker 使用异步刷盘方式以获得更好的性能。

Broker 日志管理

Broker 的默认日志路径位于 ${user.home}/logs/rocketmqlogs/。您可以通过编辑二进制包 conf 目录下的 xx.logback.xml 文件来修改日志级别和路径。

注意:请确保您的日志已妥善保护,以防止敏感信息泄露。