Day11 RocketMQ集群部署客户端使用
-
集群部署(由于服务器资源限制2G,选择部署三主)
我搭建步骤是对的,但是因为太吃内存了,2G支撑不起,所以直接game over
- 2m-noslave 两主无从节点,存在单点故障【每台机器开一个主节点,勉强能跑】
- 2m-2s-sync/async 两主两从同步/异步,主节点宕机需手动切换【在每台机器上要开一主一丛,我2G内存受限,会OOM】
- dledger 具备主从切换的高可用集群。集群中的节点会基于Raft协议随机选举出一个leader【比两主两从还更重,至少需要4G内存】
- 基于Raft协议,确保多个副本之间特别数据,保证数据强一致性
- 会接管RocketMQ原生的文件写入,写入速度稍慢,会有性能影响
▼bash复制代码# 将单机部署的机器停机 kill -9 pid # 在机器 A 上 cp /opt/rocketmq/conf/broker.conf /opt/rocketmq/conf/broker-tmp.properties vim /opt/rocketmq/conf/broker-template.properties #brokerClusterName=MyRocketMQCluster #brokerId=0 #deleteWhen=04 #fileReservedTime=48 #brokerRole=ASYNC_MASTER #flushDiskType=ASYNC_FLUSH #listenPort=10911 #storePathRootDir=/opt/rocketmq/store #storePathCommitLog=/opt/rocketmq/store/commitlog # namesrvAddr 先留空,后面统一替换 # 打包mq tar -czf rocketmq.tar.gz rocketmq # 分发到BC两台机器 # 分发到机器 B scp rocketmq.tar.gz root@192.168.1.102:/opt/ # 分发到机器 C scp rocketmq.tar.gz root@192.168.1.103:/opt/ # 解压并删除原文件 cd /opt tar -xzvf rocketmq.tar.gz rm -f rocketmq.tar.gz cd rocketmq mkdir -p store/{commitlog,consumequeue,index,checkpoint,abort} mkdir -p /var/log/rabbitmq # 复制编辑配置文件,三台机器都要改 cp /opt/rocketmq/conf/broker-tmp.properties /opt/rocketmq/conf/broker-c.properties vim /opt/rocketmq/conf/broker-c.properties # brokerName=broker-b # namesrvAddr= IP1:9876;IP2:9876;IP3:9876 # 看内存情况,禁用以下配置,机器4G以上且只部署MQ服务无需操作 # enableScheduleMessageStats=false # timerWheelEnable=false # timerColdDataCheckEnable=false # enableLmq=false # enableMultiDispatch=false # transientStorePoolEnable=false # sendMessageThreadPoolNums=1 # pullMessageThreadPoolNums=1 # 把临时文件删除 rm -f /opt/rocketmq/conf/broker-tmp.properties # 启动nameserver nohup sh /opt/rocketmq/bin/mqnamesrv > /var/log/rocketmq/namesrv.log 2>&1 & # 分别启动broker nohup /opt/rocketmq/bin/mqbroker -c /opt/rocketmq/conf/broker-a.properties > /var/log/rocketmq/broker.log 2>&1 & nohup /opt/rocketmq/bin/mqbroker -c /opt/rocketmq/conf/broker-b.properties > /var/log/rocketmq/broker.log 2>&1 & nohup /opt/rocketmq/bin/mqbroker -c /opt/rocketmq/conf/broker-c.properties > /var/log/rocketmq/broker.log 2>&1 & # PS:不出意外的话得出意外了,我2G的内存不够,启动不了,即使禁用了一些配置,AI推荐降级到4.9,但我已经乏了 # 验证集群 # 查看集群 /opt/rocketmq/bin/mqadmin clusterList -n 192.168.1.101:9876 # 应输出: # MyRocketMQCluster broker-a 192.168.1.101:10911 # MyRocketMQCluster broker-b 192.168.1.102:10911 # MyRocketMQCluster broker-c 192.168.1.103:10911 # 测试发送 export NAMESRV_ADDR=192.168.1.101:9876 sh /opt/rocketmq/bin/tools.sh org.apache.rocketmq.example.quickstart.Producer

-
运行时架构图(核心:nameServer、broker、client)
nameServer:服务注册协调中心,独立启动
broker:提供消息存储、传递、查询等功能,是RocketMQ中最繁琐也是最娇贵的组建
client:客户端
💡为什么RocketMQ不使用Zookeeper作为注册中心,而是自己实现nameServer?
- 设计简洁且专用 专门的场景设计,更加轻量,易部署
- 高性能 zookeeper比较重,繁杂的管理机制和强一致性带来了较多开销
- 高可用 无状态,多个nameserver之间对等,即使只有一台示例也不会影响整体运行,而zookpeer在节点之间同步要求严格
- 降低依赖性 自己实现可以降低外部系统的依赖,简化系统复杂度,同时能够自己掌控优化
- 定制化需求 需要实现特定的功能,不必受限与zookeeper的实现和接口

客户端编程模型
-
消息处理模型

-
消息确认机制
要支持互联网金融场景,消息安全必须是最高优先级保障。 消息安全有两方面要求,一方面是生产者要能确保将消息发送到Broker上,另一方面消费者要能确保从broker上获取到消息
-
消息生产端采用消息确认加重试机制保障消息正常发送到RocketMQ
三种发送消息的方式
- 单向发送SendOneWay 只管发送消息,不管成不成功。发送消息效率高,但如果发送失败则无法补救,适用于一些追求效率且允许消息丢失的业务场景。
- 同步发送send 发送完消息要等待broker回复结果{SEND_OK,FLUSH_DISK_TIMEOUT,FLUSH_SLAVE_TIMEOUT,SLAVE_NOT_AVAILABLE},保证消息一定发到broker,如果返回失败,可以进行重试,但返回失败不一定代表没有推送给下游消费者。能保障消息发送的安全性,但效率低,会阻塞当前线程。
- 异步发送send callback{onSuccess,onException} 不阻塞主线程,由异步回调处理成功或失败的情况,但是在回调结束之前,生产者主线程不能调用shutDown方法关闭主线程。
-
消费者端采用状态确认机制保证消费者一定能正常处理对应消息
CONSUME_SUCCESS | RECONSUME_LATER 当broker拿到的是RECONSUME_LATER状态时,会将当前消息放到对应消费者组的重试topic中,以免影响正常队列的运行,当重试次数大于16时,则将消息推入到死信topic,如仍需要重试,可以人工介入进行消息补救
-
消费者组可以自行指定起始消费位点
CONSUME_FROM_LAST_OFFSET | CONSUME_FROM_FIRST_OFFSET | CONSUME_FROM_TIMESTAMP 虽然提供了这么一种功能,但是消费者无法知道自己是在那里中断消费的,所以最好是建一个新的消费者组,再从起始点开始消费,也可以指定具有的时间点,consumer.setConsumerTimestamp("20251223171201");
-
-
消息处理模式
- 广播模式 一个消息,推送到所有消费者实例,不关心消费者组
consumer.setMessageModel(MessageModel.BROADCASTING);- 将offset交由consumer自行保管,只要consumer拉取,就返回对应消息,但是broker端不会对消费失败的消息进行重试
- 集群模式(默认) 一个消息只会推送到指定的消费者组,由消费者组中的多个实例共同消费
- 给每个ConsumerGroup维护一个统一的offset,同一个组内只会被消费一次
- 广播模式 一个消息,推送到所有消费者实例,不关心消费者组
-
消息过滤
- 简单过滤 订阅时使用tag过滤自己感兴趣的内容
- consumer.subscribe("TagFilterTest", "Tag");
- 使用两个竖线(||)连接多个Tag值,也可以使用星号(*)匹配所有
- SQL过滤 使用标准SQL语句进行过滤 WHERE 后面的过滤条件
- 数值比较,比如:>,>=,<,<=,BETWEEN…AND,=;
- 字符比较,比如:=,<>,IN;
- IS NULL 或者 IS NOT NULL;
- 逻辑符号 AND,OR,NOT;
- 顺序消息机制
- 生产者将一组有序的消息发送至同一个MessageQueue,天然保证有序性
- 消费者实现MessageListenerOrderly接口,按序拿完这一批数据后才去读其他MessageQueue的消息
- 注意事项
- 有序性一般是指局部有序,如果需要全局有序,那就只能使用一个queue,性能就很差了
- 生产者尽可能将有序消息分散到不同的queue,避免数据过于集中导致热点竞争
- 消费者端只进行有限次数的重试,如果消息处理失败,要保持有序性,就得阻塞一直重试失败消息的消费,但如果一直失败,此时就无法消费其他消息了。
- 消费者如果处理出现异常,建议返回ConsumeOrderlyStatus.SUSPEND_CURRENT_QUEUE_A_MOMENT挂起替代抛出异常
- 简单过滤 订阅时使用tag过滤自己感兴趣的内容




