### 课堂目标 1. 理解什么是MQ 与特点 2. MQ的基本原理 3. MQ使用场景 4. MQ的安装部署 # 1、什么是MQ? https://rocketmq.apache.org/ https://rocketmq.apache.org/zh/docs/ Message Queue JMS-\>rocketmq - 一个队列(先进先出) - 可以持久化 !\[image-20250831130449276\](https://woniumd.oss-cn-hangzhou.aliyuncs.com/daiwei/java/20250912085838874.png) !\[image-20260121152044995\](https://woniumd.oss-cn-hangzhou.aliyuncs.com/daiwei/java/20260121152045043.png) # 2、有什么特点?(重点) \*\*异步、解耦、削峰\*\* - \*\*异步\*\* 发送方(生产者)发出消息后,不需要等待接收方(消费者)立即处理并返回结果,而是可以立刻继续执行后续代码。接收方会在自己准备好的时候再去处理这个消息。 外卖员把外卖直接送到你的手里叫同步。 外卖员把外卖送到你们门口,发个短信或者打个电话给你,让你开门自己拿,叫异步。 系统中的场景: 数据导出: !\[image-20251211095403081\](https://woniumd.oss-cn-hangzhou.aliyuncs.com/daiwei/java/20251211174517992.png) - \*\*解耦\*\* 通过中间件(如RocketMQ),将系统间直接的、紧密的调用关系转变为间接的、松散的关系。 业务发起的服务,不需要知道其他系统的接口地址,也不需要关心它们是否在线、性能如何,它只需要按照固定格式向MQ发送消息即可。 支付系统、库存系统、短信系统也同样只需要从MQ取消息处理。 !\[image-20251211100002661\](https://woniumd.oss-cn-hangzhou.aliyuncs.com/daiwei/java/20251211174537635.png) 或者比如:有一个支付订单的业务: 1、支付订单交易成功 2、给用户发短信、或加消费积分 或者 发送优惠卷 正常情况: 1、支付、下发积分、支付成功短信会在一个service里调用,甚至是在一个事务里面。 2、如果这3个步骤,如果有一个步骤出错,整个业务就会回滚。 3、作为平台,最关心的就是订单支付成功,能否发送短信或者加积分问题,相对来说没有这么重要。 通过解耦来实现: 1、下单支付成功后,代表已经完成了核心业务。 2、下单支付完成后就分别发送两条消息(积分增加、短信通知) 3、积分增加和短信通知所对应的功能,分别监听这两类消息; 4、一旦收到消息,再去执行积分增加和短信发送的功能。 小单支付 和 积分增加、短信通知 完全隔离,这种方式成为解耦。 书面解释:通过中间件(如RocketMQ),将系统间直接的、紧密的调用关系转变为间接的、松散的关系。 业务发起的服务,不需要知道其他系统的接口地址,也不需要关心它们是否在线、性能如何,它只需要按照固定格式向MQ发送消息即可。 - \*\*削峰\*\* !\[image-20250831134235991\](https://woniumd.oss-cn-hangzhou.aliyuncs.com/daiwei/java/20251211174638058.png) 用消息队列作为\*\*"水库"或"缓冲池"\*\*,拦截并暂存上游系统突然涌来的巨大流量(请求),然后让下游系统按照自己能够承受的处理速度来消费这些消息,从而避免下游系统被突发流量冲垮。 !\[image-20260121152955029\](https://woniumd.oss-cn-hangzhou.aliyuncs.com/daiwei/java/20260121152955072.png) # 3、有哪些 消息队列? \*\*RocketMQ\*\*、RabbitMQ、Kafka、Redis... \| \| RocketMQ \| Kafka \| RabitMQ \| \| ------ \| ----------------------- \| ------------------------------ \| -------------------- \| \| 吞吐量 \| 10W+ \| 15w+ \| 6w左右 \| \| 时效性 \| 毫秒 \| 毫秒 \| 微秒 \| \| 可用性 \| 非常高,分布式架构高可用 \| 非常高,分布式架构高可用 \| 自己做主从 \| \| 可靠性 \| 0丢失 \| 0丢失 \| 高 \| \| 特点 \| 高吞吐、高可靠、低延迟 \| 超高吞吐量、效率高、\*\*零拷贝\*\* \| 功能丰富,吞吐量一般 \| 零拷贝:减少了内存之间的拷贝次数,直接从磁盘-\>内核缓冲区-\>用户缓冲区-\>socket缓冲区-\>网卡缓冲区 ,反复切换上下文,耗时 直接从磁盘-\>内核缓冲区-\>socket缓冲区-\>网卡缓冲区,并非0次,不需要从内存状-\>用户缓冲区走而已 # 4、RocketMQ的组成和原理(重点) - \*\*NameServer\*\* 多个 Broker的注册中心,负责注册发现Boroker,给生产者和消费者提供哪些broker可用。 如果有多个NameServer,每一个NameServer之间是没有通信的,在Broker上配置多个NameServer即可。 - \*\*Broker\*\* Broker是就是一个实例,就是一个MQ - \*\*Producer\*\* 创建或发布消息的生产者 - ProducerGroup 多个同类型的生产者 - \*\*Consumer\*\* 接收或订阅的消费者 - ConsumerGroup 多个同类型的消费者 - Topic(msg .email) 主题,消息存放的逻辑单元,(站内信、短信发送、支付成功通知) - Queue 消息存放的实际容器,消息存储的物理单元 !\[image-20251022103520738\](https://woniumd.oss-cn-hangzhou.aliyuncs.com/daiwei/java/20251211174713337.png) !\[微信图片_2025-08-10_113110_751\](https://woniumd.oss-cn-hangzhou.aliyuncs.com/daiwei/java/20250912085839804.png) \*\*1、有四个大的部分 NameSrv、broker、producer、consumer\*\* 2、一个Topic对用多个queue,提升效率;Topic只是逻辑存在,存数据的是 Queue 3、数据信息存放再commit log里、queue message只存放 Offset \*\*4、Master和 slave的复制规则\*\* 1. \*\*异步复制 (\`ASYNC_MASTER\`)\*\* 这是 RocketMQ \*\*默认\*\*和\*\*推荐\*\*的模式,在性能和可靠性之间取得了很好的平衡。 - \*\*工作原理\*\*: 1. 生产者将消息发送到 Master。 2. Master 将消息\*\*写入本地磁盘\*\*(根据 \`flushDiskType\` 配置决定是同步刷盘还是异步刷盘)。 3. 一旦消息成功写入 Master 的磁盘,Master \*\*立即返回"发送成功"的响应\*\*给生产者。 4. \*\*在这之后\*\*,Master 会\*\*异步地\*\*、\*\*分批地\*\*将消息数据复制给其所有的 Slave 节点。 5. Slave 节点收到数据后,再将其写入自己的磁盘。 - \*\*优点\*\*: - \*\*低延迟、高吞吐\*\*:因为生产者不需要等待数据同步到 Slave 的过程,所以写入延迟非常低,整体吞吐量很高。 - \*\*缺点\*\*: - \*\*存在少量数据丢失的风险\*\*:如果 Master 在成功响应生产者之后、但还未将消息复制到 Slave 时突然宕机且\*\*无法恢复\*\*(例如磁盘损坏),那么这条消息就会丢失,因为 Slave 上没有它的副本。 - 这种模式保证的是 \*\*RocketMQ 节点本身不丢消息\*\*(因为消息已持久化到 Master 磁盘),但无法应对主节点彻底损毁的极端情况。 2. \*\*同步复制 (\`SYNC_MASTER\`)\*\* 这种模式以牺牲部分性能为代价,换取最高的数据可靠性。 - \*\*工作原理\*\*: 1. 生产者将消息发送到 Master。 2. Master 将消息写入本地磁盘。 3. Master \*\*等待\*\*,直到消息被成功复制到\*\*至少一个\*\* Slave 节点,并且 Slave 也将其写入磁盘。 4. 收到 Slave 的成功确认后,Master \*\*才返回"发送成功"的响应\*\*给生产者。 - \*\*优点\*\*: - \*\*数据高可靠\*\*:只要主从节点不是同时宕机,消息就不会丢失。即使 Master 磁盘彻底损坏,由于消息已经在 Slave 上存在副本,数据仍然是安全的。 - \*\*缺点\*\*: - \*\*更高的延迟、更低的吞吐\*\*:因为每次写入都需要等待跨网络的复制操作完成,所以写入延迟会显著增加,整体吞吐量也会下降。 # 5、RocketMQ消息发送接收 \*\*生产者发送类型三种\*\* - \*\*同步发送\*\*:同步发送会\*\*阻塞\*\*当前主线程,直到收到 RocketMQ Broker 的确认(或超时)。这能确保你知道消息是否成功送达,但会牺牲一些性能。 syncSend()方法代表就是发送同步消息。 \`\`\` SendResult result = template.syncSend("TOPIC名称","发送的内容"); \`\`\` - \*\*异步发送:\*\*默认以异步方式发送消息。这意味着 \`send\` 方法调用后会立即返回,通过回调函数进行确认发送状态。 \`\`\` public void asyncSend() { template.asyncSend("daiwei_topic", "这是一个异步消息...", new SendCallback() { @Override public void onSuccess(SendResult sendResult) { System.out.println("=====异步消息发送成功===="+sendResult.toString()); } @Override public void onException(Throwable throwable) { System.out.println("=====异步消息发送失败===="); } }); ..... ..... } \`\`\` - \*\*单向发送:\*\*只负责发送,不管消息是否发送成功。 没有返回值,直接发送就结束,无论是否发送成功! \`\`\` /\*\* \* 发送一个单向消息 \*/ public void sendOneWay() { template.sendOneWay("daiwei_topic", "这是一个单向消息..."); } \`\`\` \*\*消费者消费消息分两种\*\*: - \*\*拉模式:\*\*消费者主动去 Broker 上拉取消息(几乎不用)。 - \*\*推模式:\*\*消费者等待 Broker 把消息推送过来 只需要实现 RocketMQListener接口的onMessage()方法即可。 onMessage(String s) 方法里面的 string就是接受的消息内容。 \`\`\` @Component @RocketMQMessageListener(topic = "daiwei_topic", consumerGroup = "daiwei_consumer_group", selectorExpression = "\*", messageModel = MessageModel.CLUSTERING, consumeMode = ConsumeMode.ORDERLY ) public class MqAaaConsumer implements RocketMQListener { @Override public void onMessage(String s) { System.out.println("消费者消费aaaaaa=========="+s); } } \`\`\` 中间件的使用方式一般都有下列几种: - 原生命令行 - 原生中间件官网的api调用 - \*\*spring boot\*\* - spring cloud alibaba - 个人开源的包装api # 6、下载安装Rocket MQ https://rocketmq.apache.org/zh/docs/ https://rocketmq.apache.org/zh/download ## 6.1选择版本\*\*5.2.0\*\* !\[image-20250809131038735\](https://woniumd.oss-cn-hangzhou.aliyuncs.com/daiwei/java/20260204103722958.png) ## 6.2修改环境变量 配置一下环境变量 ROCKETMQ_HOME \`\`\` D:\\My_WoNiu\\rocketmq-all-5.2.0-bin-release\\bin \`\`\` !\[image-20250831194924367\](https://woniumd.oss-cn-hangzhou.aliyuncs.com/daiwei/java/20260204103722916.png) ## 6.3修改内存大小配置 !\[image-20250831195024913\](https://woniumd.oss-cn-hangzhou.aliyuncs.com/daiwei/java/20260204103722935.png) 修改内存大小 !\[image-20250831195131843\](https://woniumd.oss-cn-hangzhou.aliyuncs.com/daiwei/java/20260204103722926.png) !\[image-20250831195105828\](https://woniumd.oss-cn-hangzhou.aliyuncs.com/daiwei/java/20260204103722998.png) ## 6.4启动nameserver \*\*管理生产者和消费者;当作MQ的注册中心\*\* \`\`\` start mqnamesrv.cmd \`\`\` !\[image-20250831195324206\](https://woniumd.oss-cn-hangzhou.aliyuncs.com/daiwei/java/20260204103722951.png) \`\`\` jps查看是否启动成功 \`\`\` !\[image-20250831195239206\](https://woniumd.oss-cn-hangzhou.aliyuncs.com/daiwei/java/20260204103723582.png) ## 6.5修改broker配置 !\[image-20250831195600123\](https://woniumd.oss-cn-hangzhou.aliyuncs.com/daiwei/java/20260204103723660.png) 开启自动创建topic(开发和测试环境) 连接道NameServer \`\`\` autoCreateTopicEnable=true namesrvAddr=localhost:9876 \`\`\` - Broker 有很多重要的配置项 - \`brokerClusterName\`(所属集群名) - \`brokerName\`(Broker名称) - \`brokerId\`(0 表示 Master,\>0 表示 Slave) - \`deleteWhen\`(何时删除过期文件) - \`fileReservedTime\`(文件保留时间) - \`brokerRole\`(角色,ASYNC_MASTER 等) - \`flushDiskType\`(刷盘方式,ASYNC_FLUSH 等) !\[image-20250831203001627\](https://woniumd.oss-cn-hangzhou.aliyuncs.com/daiwei/java/20260204103723699.png) ## 6.6启动broker \*\*broker管理消息,把消息提供给消费者进行消费\*\* \`\`\` mqbroker.cmd -n localhost:9876 -c E:\\APP\\rocketmq\\rocketmq-all-5.2.0-bin-release\\conf\\broker.conf \`\`\` !\[image-20250809132133297\](https://woniumd.oss-cn-hangzhou.aliyuncs.com/daiwei/java/20260204103723709.png) 启动完成 !\[image-20250831203340835\](https://woniumd.oss-cn-hangzhou.aliyuncs.com/daiwei/java/20260204103723727.png) jps看一下 !\[image-20250831203304146\](https://woniumd.oss-cn-hangzhou.aliyuncs.com/daiwei/java/20260204103723704.png) # 7、自带的命令工具 mqadmin.cmd(可选) ## 7.1创建topic(updatetopic) \*\*\*\* \`\`\` mqadmin.cmd updatetopic -n localhost:9876 -t MyTestTopic -c DefaultCluster \`\`\` 返回如下信息代表成功! \`\`\` D:\\My_WoNiu\\rocketmq-all-5.2.0-bin-release\\bin\>mqadmin.cmd updatetopic -n localhost:9876 -t MyTestTopic -c DefaultCluster create topic to 192.168.0.59:10911 success. TopicConfig \[topicName=MyTestTopic, readQueueNums=8, writeQueueNums=8, perm=RW-, topicFilterType=SINGLE_TAG, topicSysFlag=0, order=false, attributes={}\] \`\`\` ## 7.2 查看topic(topicList) \`\`\` mqadmin topicList -n localhost:9876 \`\`\` ## 7.3 查看有消费者组 \`\`\` mqadmin consumerProgress -n localhost:9876 \`\`\` ## 7.4 查看topic信息 \`\`\` mqadmin topicstatus -n localhost:9876 -t topic126 \`\`\` ## 7.5 发送消息 \`\`\` mqadmin sendMessage -n {NameServer地址} -t {Topic名称} -p "消息体" mqadmin sendMessage -n "192.168.1.8:9876" -t topic-test -p "hello test" \`\`\` ## 7.6模拟消费者 \`\`\` mqadmin consumeMessage -n {NameServer地址} -t {Topic名称} -c {消费者组名称} mqadmin consumeMessage -n "192.168.1.100:9876" -t topic-test -c my_test_consumer_group \`\`\` ## 7.7 其他mqadmin命令\*\*(可选)\*\* \`\`\` mqadmin 工具的命令很丰富,按功能可以分为以下几大类。所有命令都需要 -n 参数指定 NameServer 地址,且都可以通过 -h 参数获取详细帮助。 主题(Topic)管理 updateTopic / createTopic:创建或更新主题配置。关键参数:-t 主题名、-c 集群名、-r/-w 读写队列数(默认8)。 deleteTopic:删除主题。关键参数:-t 主题名、-c 集群名。 topicList:查看主题列表。可选参数:-c 可查看主题的集群和订阅信息。 topicRoute:查看主题路由信息(队列分布等)。关键参数:-t 主题名。 topicStatus:查看主题消息队列的偏移量(Offset)。关键参数:-t 主题名。 topicClusterList:查看主题所在集群列表。关键参数:-t 主题名。 updateTopicPerm:更新主题的读写权限。关键参数:-t 主题名、-p 权限(W=2, R=4, WR=6)。 statsAll:打印主题的订阅关系、TPS、积压量等统计信息。 集群与Broker管理 clusterList:查看集群信息(Broker名称、ID、TPS等)。可选参数:-m 打印更多信息。 brokerStatus:获取 Broker 运行状态。关键参数:-b Broker地址 或 -c 集群名。 brokerConsumeStats:获取Broker消费统计数据(含积压量Diff)。关键参数:-b Broker地址。 cleanUnusedTopic:清理未使用的主题。关键参数:-b Broker地址 或 -c 集群名。 deleteExpiredCommitLog:删除过期的CommitLog文件以释放磁盘空间。关键参数:-n NameServer地址。 cleanExpiredCQ:清理过期的消费队列文件。关键参数:-b Broker地址 或 -c 集群名。 消费者与消费组管理 updateSubGroup:创建或更新消费者的订阅组配置。关键参数:-g 消费组名。 deleteSubGroup:删除订阅组。关键参数:-g 消费组名。 \`\`\` # 8、安装Rocket MQ DashBoard https://rocketmq.apache.org/zh/download 往下拉到底 !\[image-20251021173547123\](https://woniumd.oss-cn-hangzhou.aliyuncs.com/daiwei/java/20260204103724162.png) 1、看看ReadMe !\[image-20251021173848554\](https://woniumd.oss-cn-hangzhou.aliyuncs.com/daiwei/java/20260204103724305.png) 2、直接运行(不建议) \`\`\` mvn spring-boot:run \`\`\` 3、打包成jar再运行 - 进入工程的目录 !\[\](https://woniumd.oss-cn-hangzhou.aliyuncs.com/daiwei/java/20260204103724274.png) - 运行maven打包脚本 \`\`\` mvn clean package -Dmaven.test.skip=true \`\`\` !\[image-20251021174924690\](https://woniumd.oss-cn-hangzhou.aliyuncs.com/daiwei/java/20260204103724425.png) - 运行jar \`\`\` java -jar target/rocketmq-dashboard-1.0.1-SNAPSHOT.jar \`\`\` !\[image-20251021175110665\](https://woniumd.oss-cn-hangzhou.aliyuncs.com/daiwei/java/20260204103724412.png) 打开控制台 \`\`\` http://localhost:8080 \`\`\` !\[image-20260204114451862\](https://woniumd.oss-cn-hangzhou.aliyuncs.com/daiwei/java/20260204114452003.png) # 9、DockerCompose安装(可选) \`\`\` version: '3.5' services: rmqnamesrv: image: foxiswho/rocketmq:server container_name: rmqnamesrv ports: - 9876:9876 volumes: - ./data/logs:/opt/logs - ./data/store:/opt/store networks: rmq: aliases: - rmqnamesrv rmqbroker: image: foxiswho/rocketmq:broker container_name: rmqbroker ports: - 10909:10909 - 10911:10911 volumes: - ./data/logs:/opt/logs - ./data/store:/opt/store - ./data/broker.conf:/etc/rocketmq/broker.conf environment: NAMESRV_ADDR: "rmqnamesrv:9876" JAVA_OPTS: " -Duser.home=/opt" JAVA_OPT_EXT: "-server -Xms128m -Xmx128m -Xmn128m" command: mqbroker -c /etc/rocketmq/broker.conf depends_on: - rmqnamesrv networks: rmq: aliases: - rmqbroker rmqconsole: image: styletang/rocketmq-console-ng container_name: rmqconsole ports: - 18082:8080 environment: JAVA_OPTS: "-Drocketmq.namesrv.addr=rmqnamesrv:9876 -Dcom.rocketmq.sendMessageWithVIPChannel=false" depends_on: - rmqnamesrv networks: rmq: aliases: - rmqconsole networks: rmq: name: rmq driver: bridge \`\`\`
20-3阶内容-2.2.1.0-MQ安装
https://xiaochenblog.icu/archives/b3d56ce1-98bd-466a-8d13-5943ec4c604d
评论