4、RocketMQ最佳实践

Navy2026-09-10mqrocketmq

所谓最佳实践,主要针对的是在实际生产环境中可能会涉及到的需要调优的部分,综合考虑主要有六大环节:producer、broker、consumer、nameserver、jvm以及os,接下来就从这六个方面总结可能会用到的一些配置

1、Producer

① 、一个应用尽可能用一个topic,消息的子类型用tags来标识

rocketmq支持应用给消息自由设置tags,消费者订阅主题的时候可以通过匹配规则指定要消费带有哪些tags的消息,从而实现消息的过滤,比如:

// 消费者订阅主题
consumer.subscribe("TopicTest", "*");

// 生产者发送消息
Message message = new Message("TopicTest", "TagA", ("Hello MQ").getBytes(RemotingHelper.DEFAULT_CHARSET));

②、每个消息在业务层面的唯一标识码要应用于keys字段

服务器会为每个消息建立一个hash索引,应用可以通过topic、keys来查询这条消息的内容,甚至是trace,这样就非常有利于后续出现问题的时候的问题定位。但是因为是哈希索引,所以务必需要保证key的唯一性,可以有效减少哈希冲突

③、消息不管是发送成功还是发送失败,都需要打印完整的日志(SendResult和Key字段)

④、如果相同性质的消息量比较大,建议可以使用批量消息,在一定程度上可以有效提升性能

⑤、消息大小建议不要超过512kb

⑥、同步send发送消息会阻塞,如果有性能要求时,可以使用异步发送的方式:send(msg,callback)

⑦、如果在一个jvm中有多个生产者进行大数据处理,可以用少数生产者(3~5个)使用异步发送方式就可以了,通过setInstanceName给每个生产者设置一个示例名称

⑧、send发送消息只要不抛出异常,就表示发送成功。

但是发送成功有多种状态,在SendResult里定义:

  • SEND_OK:消息发送成功
  • FLUSH_DISK_TIMEOUT:消息发送成功,但是服务器刷盘的时候超时了,消息已经进入消息队列,只有此时服务器宕机消息才会丢失
  • FLUSH_SLAVE_TIMEOUT:消息发送成功,但是消息同步到slave超时,消息已经进入到队列,只有此时服务器宕机消息才会丢失
  • SLAVE_NOT_AVAILABLE:消息发送成功,但是此时slave不可用,消息已经进入队列,只有此时服务器宕机消息才会丢失

如果状态是FLUSH DISK TIMEOUT或FLUSH SLAVE TIMEOUT,并且broker正好关闭,可以丢弃这条消息,或者重发(建议重发,消费端自己去重)。

Producer 向Broker 发送请求会等待响应,但如果达到最大等待时间,未得到响应,则客户端将抛出RemotingTimeoutException。默认等待时间是3秒,如果使用send(msg, timeout),则可以自己设定超时时间,但超时时间不能设置太小,因为Borker 需要一些时间来刷新磁盘或与从属设备同步。如果该值超过syncFlushTimeout,则该值可能影响不大,因为Broker可能会在超时之前返回FLUSH_SLAVE_TIMEOUT或FLUSH_SLAVE_TIMEOUT的响应

⑨、对于消息不可丢失应用,务必要有消息重发机制

Producer的send方法本身支持内部重试,至多重试3次,如果发送失败,则轮转到下一个Broker。 这个方法的总耗时时间不超过sendMsgTimeout设置的值,默认10s。

所以,如果本身向broker发送消息产生超时异常,就不会再做重试,为了保证消息一定发送成功,可以先将消息写入DB,并带有发送状态,只要不是发送成功的消息就由后台线程一直定时重试,直到发送成功,然后改写数据库对应消息的发送状态为成功,从而保证消息一定到达broker

2、broker

①、集群部署角色选择

Broker 角色分为ASYNC_MASTER(异步主机)、SYNC_MASTER(同步主机)以及SLAVE(从机)。

  • 如果对消息的可靠性要求比较严格,可以采用 SYNC_MASTER加SLAVE的部署方式。

  • 如果对消息可靠性要求不高,可以采用ASYNC_MASTER加SLAVE的部署方式。

  • 如果只是测试方便,则可以选择仅ASYNC_MASTER或仅SYNC_MASTER的部署方式即可

②、刷盘方式

SYNC_FLUSH(同步刷新)相比于ASYNC_FLUSH(异步处理)会损失很多性能,但是也更可靠,所以需要根据实际的业务场景做好权衡

③、配置

参数名默认值说明
listenPort10911接受客户端连接的监听端口
namesrvAddrnullnameServer 地址
brokerIP1网卡的 InetAddress当前 broker 监听的 IP
brokerIP2跟 brokerIP1 一样存在主从 broker 时,如果在 broker 主节点上配置了 brokerIP2 属性,broker 从节点会连接主节点配置的 brokerIP2 进行同步
brokerNamenullbroker 的名称
brokerClusterNameDefaultCluster本 broker 所属的 Cluster 名称
brokerId0broker id, 0 表示 master, 其他的正整数表示 slave
storePathCommitLog$HOME/store/commitlog/存储 commit log 的路径
storePathConsumerQueue$HOME/store/consumequeue/存储 consume queue 的路径
mappedFileSizeCommitLog1024 * 1024 * 1024(1G)commit log 的映射文件大小
deleteWhen04在每天的什么时间删除已经超过文件保留时间的 commit log
fileReservedTime72以小时计算的文件保留时间
brokerRoleASYNC_MASTERSYNC_MASTER/ASYNC_MASTER/SLAVE
flushDiskTypeASYNC_FLUSHSYNC_FLUSH/ASYNC_FLUSH SYNC_FLUSH 模式下的 broker 保证在收到确认生产者之前将消息刷盘。ASYNC_FLUSH 模式下的 broker 则利用刷盘一组消息的模式,可以取得更好的性能。

④、日志

Broker 的默认日志路径在 ${user.home}/logs/rocketmqlogs/ 下,可以通过修改二进制包中 conf 文件夹下的 xx.logback.xml 文件来进行日志级别和路径的修改

注意请保管好您的日志,以免敏感信息发生泄漏。

⑤、监控

一般在生产环境下各个中间件都会配置完备的监控方案,可以采用promethus进行指标数据采集,然后通过Graphana面板进行展示和告警

3、consumer

①、消费组和订阅

不同的消费群体可以独立地消费同样的主题,并且每个消费者都有自己的消费偏移量(offsets),确保同一组中的每个消费者订阅相同的主题

②、顺序消费(MessageListenerOrderly)

消费者将锁定每个MessageQueue,以确保每个消息被一个按顺序使用,这将导致性能损失,但如果关心消息的顺序时,它就很有用了。不建议抛出异常,可以返回ConsumeOrderlyStatus.SUSPEND_CURRENT_QUEUE_A_MOMENT代替

③、并行消费(MessageListenerConcurrently)

顾名思义,消费者将同时使用这些消息。为良好的性能,推荐使用此方法。不建议抛出异常,可以返回ConsumeConcurrentlyStatus.RECONSUME_LATER代替。

③、消费状况(Consume Status)

对于 MessageListenerConcurrently,可以返回 RECONSUME_LATER 告诉消费者,当前不能消费它并且希望以后重新消费。然后可以继续使用其他消息。

对于MessageListenerOrderly,如果关心顺序,就不能跳过消息,可以返回 SUSPEND_CURRENT_QUEUE_A_MOMENT 来告诉消费者等待片刻。

④、阻塞

不建议阻塞 Listener,因为它会阻塞线程池,最终可能会停止消费程序。

⑤、线程数

消费者使用一个 ThreadPoolExecutor 来处理内部的消费,因此可以通过设置setConsumeThreadMin或setConsumeThreadMax来更改它。

⑥、从何处开始消费

当建立一个新的 Consumer Group 时,需要决定是否需要消费 Broker 中已经存在的历史消息。rocketmq提供了下面几个选择:

  • CONSUME_FROM_LAST_OFFSET:将忽略历史消息,并消费此后生成的任何内容。

  • CONSUME_FROM_FIRST_OFFSET:将消耗 Broker 中存在的所有消息。

  • CONSUME_FROM_TIMESTAMP:消费在指定的时间戳之后生成的消息。

前两个跟kafka非常类似,指定时间戳这个稍有不同

⑦、消息幂等性

RocketMQ无法避免消息重复,如果业务对重复消费非常敏感,务必在业务层面做幂等去重:

  • 通过记录消息唯一键进行去重
  • 使用业务层面的状态机制去重

4、nameserver

在Apache RocketMQ中,NameServer 用于协调分布式系统的每个组件,主要通过管理主题路由信息来实现协调。管理由两部分组成:

  • Brokers 定期更新保存在每个名称服务器中的元数据。

  • 名称服务器是为客户端提供最新的路由信息服务的,包括生产者、消费者和命令行客户端。

因此,在启动 brokers 和 clients 之前,我们需要告诉他们如何通过给他们提供的一个名称服务器地址列表来访问名称服务器。在Apache RocketMQ中,可以用四

种方式完成:

①、编程方式

对于 brokers,我们可以在 broker 的配置文件中指定:

namesrvAddr=name-server-ip1:port;name-server-ip2:port

对于生产者和消费者,我们可以给他们提供姓名服务器地址列表如下:

String NAME_SERVER_ADDR = "192.168.0.61:9876;192.168.0.62:9876;192.168.0.63:9876";
// 消费者指定nameserverr
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("GROUP_TEST");
consumer.setNamesrvAddr(NAME_SERVER_ADDR);

//生产者指定nameserver
DefaultMQProducer producer = new DefaultMQProducer("GROUP_TEST");
producer.setNamesrvAddr(NAME_SERVER_ADDR);

②、Java参数(参数名称不能写错)

NameServer 的地址列表也可以通过java 参数rocketmq.namesrv.addr,在启动之前指定

③、环境变量(环境变量名称不能写错)

可以设置NAMESRV_ADDR环境变量。如果设置了,Broker和clients将检查并使用其值

④、HTTP端点(HTTP Endpoint)

如果没有使用前面提到的方法指定NameServer地址列表,Apache RocketMQ将每2分钟发送一次HTTP请求,以获取和更新NameServer地址列表,初始延迟10秒。默认情况下,访问的HTTP地址是:http://jmenv.tbsite.net:8080/rocketmq/nsaddr

通过Java参数rocketmq.namesrv.domain,可以修改jmenv.tbsite.net

通过Java参数rocketmq.namesrv.domain.subgroup,可以修改nsaddr

这四种方式的优先级:

编程方式 > Java 参数 > 环境变量 > HTTP方式

5、jvm

①、推荐使用最新发布的JDK(>=1.8版本),使用服务器编译器和8g堆。设置相同的Xms和Xmx值,以防止JVM调整堆大小以获得更好的性能。简单的JVM配置:-server -Xms8g -Xmx8g -Xmn4g

②、如果不关心Broker的启动时间,可以预先触摸Java堆,以确保在JVM初始化期间分配页是更好的选择:-XX:+AlwaysPreTouch

③、禁用偏置锁定可能会减少JVM暂停:-XX:-UseBiasedLocking

④、对于垃圾回收,建议使用G1收集器:-XX:+UseG1GC -XX:G1HeapRegionSize=16m -XX:G1ReservePercent=25 -XX:InitiatingHeapOccupancyPercent=30

这些GC选项看起来有点激进,但事实证明它在生产环境中具有良好的性能,GCPauseMillis不要设置太小的值,否则JVM将使用一个小的新生代,这将导致非常频繁的新生代GC。

⑤、推荐使用滚动GC日志文件:-XX:+UseGCLogFileRotation -XX:NumberOfGCLogFiles=5 -XX:GCLogFileSize=30m

⑥、如果写入GC文件会增加代理的延迟,请将重定向GC日志文件考虑在内存文件系统中:-Xloggc:/dev/shm/mq_gc_%p.log

6、os

在bin目录中,有一个os.sh脚本列出了许多内核参数,只需要稍微的修改,就可以用于生产环境:

image-20260910113624615

以下参数需要注意(详细信息请参考:https://www.kernel.org/doc/Documentation/sysctl/vm.txt):

①、vm.extra_free_kbytes

告诉虚拟机在启动后台回收(kswapd)的阈值之间保留额外的可用内存。RocketMQ 使用此参数来避免内存分配中的高延迟。

②、vm.min_free_kbytes

如果将其设置为低于1024KB,系统将会被微妙地破坏,并且在高负载下容易出现死锁

③、vm.max_map_count

限制进程可能拥有的内存映射区域的最大数量。RocketMQ将使用mmap来加载CommitLog和ConsumeQueue,因此建议为此参数设置一个较大的值。

④、vm.swappiness

定义内核如何有效地交换内存页面。值越大,交换量越大,值越小,交换量越小。

⑤、File descriptor limits

RocketMQ 需要文件(CommitLog和ConsumeQueue)和网络连接的打开文件描述符。建议将文件描述符限制为655350 。

⑥、Disk scheduler

RocketMQ建议使用截止时间I/O调度程序,它试图为请求提供有保证的延迟

参考资料:

https://access.redhat.com/documentation/en-US/Red_Hat_Enterprise_Linux/6/html/Performance_Tuning_Guide/ch06s04s02.html


扫码关注

最后更新时间: 9/19/2026, 5:28:30 AM