当前位置:

RabbitMQ

故人归 |  2025-07-13
 43人浏览

文档地址:https://www.kuangstudy.com/zl/rabbitmq

代码:https://gitee.com/zhayuyao/my-rabbitmq

RabbitMQ是一个开源的遵循AMQP协议实现的基于Erlang语言编写,支持多种客户端(语言)。用于在分布式系统中存储消息,转发消息,具有高可用,高可扩性,易用性等特征。

1、什么是中间件

1.1、什么是中间件

我国企业从20世纪80年代开始就逐渐进行信息化建设,由于方法和体系的不成熟,以及企业业务和市场需求的不断变化,一个企业可能同时运行着多个不同的业务系统,这些系统可能基于不同的操作系统、不同的数据库、异构的网络环境。现在的问题是,如何把这些信息系统结合成一个有机地协同工作的整体,真正实现企业跨平台、分布式应用。中间件便是解决之道,它用自己的复杂换取了企业应用的简单。

中间件(Middleware)是处于操作系统和应用程序之间的软件,也有人认为它应该属于操作系统中的一部分。人们在使用中间件时,往往是一组中间件集成在一起,构成一个平台(包括开发平台和运行平台),但在这组中间件中必须要有一个通信中间件,即中间件=平台+通信,这个定义也限定了只有用于分布式系统中才能称为中间件,同时还可以把它与支撑软件和实用软件区分开来。

举例:

  1. RMI(Remote Method Invocations, 远程调用)
  2. Load Balancing(负载均衡,将访问负荷分散到各个服务器中)
  3. Transparent Fail-over(透明的故障切换)
  4. Clustering(集群,用多个小的服务器代替大型机)
  5. Back-end-Integration(后端集成,用现有的、新开发的系统如何去集成遗留的系统)
  6. Transaction事务(全局/局部)全局事务(分布式事务)局部事务(在同一数据库联接内的事务)
  7. Dynamic Redeployment(动态重新部署,在不停止原系统的情况下,部署新的系统)
  8. System Management(系统管理)
  9. Threading(多线程处理)
  10. Message-oriented Middleware面向消息的中间件(异步的调用编程)
  11. Component Life Cycle(组件的生命周期管理)
  12. Resource pooling(资源池)
  13. Security(安全)
  14. Caching(缓存)

1.2、为什么需要中间件系统

具体地说,中间件屏蔽了底层操作系统的复杂性,使程序开发人员面对一个简单而统一的开发环境,减少程序设计的复杂性,将注意力集中在自己的业务上,不必再为程序在不同系统软件上的移植而重复工作,从而大大减少了技术上的负担。中间件带给应用系统的,不只是开发的简便、开发周期的缩短,也减少了系统的维护、运行和管理的工作量,还减少了计算机总体费用的投入。

1.3、中间件特点

为解决分布异构问题,人们提出了中间件(middleware)的概念。中间件是位于平台(硬件和操作系统)和应用之间的通用服务,针对不同的操作系统和硬件平台,它们可以有符合接口和协议规范的多种实现。

也许很难给中间件一个严格的定义,但中间件应具有如下的一些特点:

  • 满足大量应用的需要
  • 运行于多种硬件和OS平台
  • 支持分布计算,提供跨网络、硬件和OS平台的透明性的应用或服务的交互
  • 支持标准的协议
  • 支持标准的接口

由于标准接口对于可移植性和标准协议对于互操作性的重要性,中间件已成为许多标准化工作的主要部分。对于应用软件开发,中间件远比操作系统和网络服务更为重要,中间件提供的程序接口定义了一个相对稳定的高层应用环境,不管底层的计算机硬件和系统软件怎样更新换代,只要将中间件升级更新,并保持中间件对外的接口定义不变,应用软件几乎不需任何修改,从而保护了企业在应用软件开发和维护中的重大投资。

1.4、在项目中什么时候时候中间件技术

在项目的架构和重构中,使用任何技术和架构的改变我们都需要谨慎斟酌和思考,因为任何技术的融入和变化都可能人员,技术,和成本的增加,中间件的技术一般现在一些互联网公司或者项目中使用比较多,如果你仅仅还只是一个初创公司建议还是使用单体架构,最多加个缓存中间件即可,不要盲目追求新或者所谓的高性能,而追求的背后一定是业务的驱动和项目的驱动,因为一旦追求就意味着你的学习成本,公司的人员结构以及服务器成本,维护和运维的成本都会增加,所以需要谨慎选择和考虑。

但是作为一个开放人员,一定要有学习中间件技术的能力和思维,否则很容易当项目发展到一个阶段在去掌握估计或者在面试中提及,就会给自己带来不小的困扰,在当今这个时代这些技术也并不是什么新鲜的东西,如果去掌握和挖掘最关键的还是自己花时间和花精力去探讨和研究。

1.5、课程的规划和安排

  • 消息中间件 ActiveMQ
  • 消息中间件 ActiveMQ
  • 消息中间件 Kafaka
  • 消息中间件 RocketMQ
  • 消息中间件应用场景说明
  • 负载均衡中间件(Nginx/Lvs)
  • 缓存中间件(Memcache/Redis)
  • 数据库中间件(ShardingJdbc/Mycat)

2、中间件技术及架构的概述

知识图谱

11111

2.1、学习中间件的方法和技巧

  1. 理解中间件在项目架构中的作用,以及各中间件的底层实现
  2. 可以使用一些类比的生活概念去理解中间件
  3. 使用一些流程图或者脑图的方式去梳理各个中间件在架构中的作用
  4. 尝试用java技术去实现中间件的原理
  5. 静下来去思考中间件在项目中设计的和使用的原因
  6. 如果找到对应的替代总结方案
  7. 尝试编写博文总结同类中间件技术的对比和使用场景
  8. 学会查看中间件的源码以及开开源项目和博文

2.2、学习目标

  • 什么是消息中间件
  • 什么是协议
  • 什么是持久化
  • 消息分发
  • 消息的高可用
  • 消息的集群
  • 消息的容错
  • 消息的冗余

2.3、什么是消息中间件

在实际的项目中,大部分的企业项目开发中,在早期都采用的是单体的架构模式,如下图:

image-20230515213548739

2.4、单体架构

在企业开发的中,大部分的初期架构都采用的是单体架构的模式进行架构,而这种架构的典型的特点:就是把所有的业务和模块,源代码,静态资源文件等都放在一个一工程中,如果其中的一个模块升级或者迭代发生一个很小变动都会重新编译和重新部署项目。 这种的架构存在的问题就是:

  1. 耦合度太高
  2. 运维的成本过高
  3. 不易维护
  4. 服务器的成本高
  5. 以及升级架构的复杂度也会增大

这样就有后续的分布式架构系统。如下

2.5、分布式架构

image-20230515213716795

何谓分布式系统呢:

通俗一点:就是一个请求由服务器端的多个服务(服务或者系统)协同处理完成

和单体架构不同的是,单体架构是一个请求发起jvm调度线程(确切的是tomcat线程池)分配线程Thread来处理请求直到释放,而分布式系统是:一个请求是由多个系统共同来协同完成,jvm和环境都可能是独立。

分布式架构系统存在的特点和问题如下:

存在问题:

  1. 学习成本高,技术栈过多
  2. 运维成本和服务器成本增高
  3. 人员的成本也会增高
  4. 项目的负载度也会上升
  5. 面临的错误和容错性也会成倍增加
  6. 占用的服务器端口和通讯的选择的成本高
  7. 安全性的考虑和因素逼迫可能选择RMI/MQ相关的服务器端通讯

好处:

  1. 服务系统的独立,占用的服务器资源减少和占用的硬件成本减少

    确切的说是:可以合理的分配服务资源,不造成服务器资源的浪费

  2. 系统的独立维护和部署,耦合度降低,可插拔性

  3. 系统的架构和技术栈的选择可以变的灵活(而不是单纯的选择java)

  4. 弹性的部署,不会造成平台因部署造成的瘫痪和停服的状态

3、基于消息中间件的分布式系统的架构

3.1、基于消息中间件的分布式系统架构

image-20230515222331018

从上图中可以看出来,消息中间件的是

  1. 利用可靠的消息传递机制进行系统和系统直接的通讯
  2. 通过提供消息传递和消息的排队机制,它可以在分布式系统环境下扩展进程间的通讯

3.2、消息中间件应用场景

  1. 跨系统数据传递
  2. 高并发的流量削峰
  3. 数据的分发和异步处理
  4. 大数据分析与传递
  5. 分布式事务

比如你有一个数据要进行迁移或者请求并发过多的时候,比如你有10W的并发请求下订单,我们可以在这些订单入库之前,我们可以把订单请求堆积到消息队列中,让它稳健可靠的入库和执行。

3.3、常见的消息中间件

ActiveMQ、RabbitMQ、Kafka、RocketMQ等。

3.4、消息中间件的本质和设计

它是一种接受数据,接受请求、存储数据、发送数据等功能的技术服务。

MQ消息队列:负责数据的传接受,存储和传递,所以性能要过于普通服务和技术。

image-20230515222611486

谁来生产消息,存储消息和消费消息呢?

image-20230515222639421

3.5、消息中间件的核心组成部分

  1. 消息的协议
  2. 消息的持久化机制
  3. 消息的分发策略
  4. 消息的高可用,高可靠
  5. 消息的容错机制

3.6、小结

其实不论选择单体架构还是分布式架构都是项目开发的一个阶段,在什么阶段选择适合的架构方式,而不能盲目追求,最后造成的后果和问题都需要自己买单。但是作为一个开发人员学习和探讨新的技术是我们每个程序开发者都应该去保持和思考的问题。当我们没办法去改变社会和世界的时候,我们为了生活和生存那就必须要迎合企业和市场的需求,发挥你的价值和所学的才能,创造价值和实现自我。

4、消息队列协议

4.1、什么是协议

image-20230515223026489

我们知道消息中间件负责数据的传递,存储,和分发消费三个部分,数据的存储和分发的过程中肯定要遵循某种约定成俗的规范,你是采用底层的TCP/IP,UDP协议还是其他的自己取构建等,而这些约定成俗的规范就称之为:协议。

所谓协议是指:
1.计算机底层操作系统和应用程序通讯时共同遵守的一组约定,只有遵循共同的约定和规范,系统和底层操作系统之间才能相互交流
2.和一般的网络应用程序的不同它主要负责数据的接受和传递,所以性能比较的高
3.协议对数据格式和计算机之间交换数据都必须严格遵守规范

4.2、网络协议的三要素

  1. 语法

    语法是用户数据与控制信息的结构与格式,以及数据出现的顺序。

  2. 语义

    语义是解释控制信息每个部分的意义。它规定了需要发出何种控制信息,以及完成的动作与做出什么样的响应。

  3. 时序

    时序是对事件发生顺序的详细说明。

比如我MQ发送一个信息,是以什么数据格式发送到队列中,然后每个部分的含义是什么,发送完毕以后的执行的动作,以及消费者消费消息的动作,消费完毕的响应结果和反馈是什么,然后按照对应的执行顺序进行处理。如果你还是不理解:大家每天都在接触的http请求协议:

1.语法:http规定了请求报文和响应报文的格式
2.语义:客户端主动发起请求称之为请求。(这是一种定义,同时你发起的是post/get请求)
3.时序:一个请求对应一个响应。(一定先有请求在有响应,这个是时序)

而消息中间件采用的并不是http协议,而常见的消息中间中间件协议有:OpenWire、AMQP、MQTT、Kafka,OpenMessage协议

面试题:为什么消息中间件不直接使用http协议呢?
1. 因为http请求报文头和响应报文头是比较复杂的,包含了cookie,数据的加密解密,状态码,响应码等附加的功能,但是对于一个消息而言,我们并不需要这么复杂,也没有这个必要性,它其实就是负责数据传递,存储,分发就行,一定要追求的是高性能。尽量简洁,快速。
2. 大部分情况下http大部分都是短链接,在实际的交互过程中,一个请求到响应很有可能会中断,中断以后就不会就行持久化,就会造成请求的丢失。这样就不利于消息中间件的业务场景,因为消息中间件可能是一个长期的获取消息的过程,出现问题和故障要对数据或消息进行持久化等,目的是为了保证消息和数据的高可靠和稳健的运行。

4.3、AMQP协议

AMQP:(全称:Advanced Message Queuing Protocol) 是高级消息队列协议。由摩根大通集团联合其他公司共同设计。是一个提供统一消息服务的应用层标准高级消息队列协议,是应用层协议的一个开放标准,为面向消息的中间件设计。基于此协议的客户端与消息中间件可传递消息,并不受客户端/中间件不同产品,不同的开发语言等条件的限制。Erlang中的实现有RabbitMQ等。

特性:

  1. 分布式事务支持
  2. 消息的持久化支持
  3. 高性能和高可靠的消息处理优势

image-20230515224207372

4.4、MQTT协议

MQTT协议:(Message Queueing Telemetry Transport)消息队列是IBM开放的一个即时通讯协议,物联网系统架构中的重要组成部分。

特点

  1. 轻量
  2. 结构简单
  3. 传输快,不支持事务
  4. 没有持久化设计

应用场景:

  1. 使用于计算能力优先
  2. 低带宽
  3. 网络不稳当的场景

image-20230515224421122

4.5、OpenMessage协议

image-20230515224432767

是近几年由阿里、雅虎和滴滴出行、Stremalio等公司共同参与创立的分布式消息中间件、流处理等领域的应用开发标准。

特点:

  1. 结构简单
  2. 解析速度快
  3. 支持事务和持久化设计

4.6、kafka协议

image-20230515224534047

Kafka协议是基于TCP/IP的二进制协议。消息内部是通过长度来分割,由一些基本数据类型组成。

特点:

  1. 结构简单
  2. 解析速度快
  3. 无事务支持
  4. 有持久化设计

4.7、小结

协议:是在tcp/ip协议基础之上构建的一种约定成俗的规范和机制、它的主要目的可以让客户端(应用程序 java,go)进行沟通和通讯。并且这种协议下规范必须具有持久性,高可用,高可靠的性能。

5、消息队列持久化

5.1、持久化

简单来说就是将数据存入磁盘,而不是存在内存中随服务器重启断开而消失,使数据能够永久保存。

image-20230515224817626

5.2、常见的持久化方式

  ActiveMQ RabbitMQ kafka RocketMQ
文件存储 支持 支持 支持 支持
数据库 支持 / / /

6、消息的分发策略

6.1、消息的分发策略

MQ消息队列有如下几个角色

  1. 生产者
  2. 存储消息
  3. 消费者

那么生产者生成消息以后,MQ进行存储,消费者是如何获取消息的呢?一般获取数据的方式无外乎推(push)或者拉(pull)两种方式,典型的git就有推拉机制,我们发送的http请求就是一种典型的拉取数据库数据返回的过程。而消息队列MQ是一种推送的过程,而这些推机制会适用到很多的业务场景也有很多对应推机制策略。

6.2、场景分析1

image-20230515225243910

比如我在APP上下了一个订单,我们的系统和服务很多,我们如何得知这个消息被哪个系统或者哪些服务或者系统进行消费,那这个时候就需要一个分发的策略。这就需要消费策略。或者称之为消费的方法论。

6.3、场景分析2

image-20230515225407511

在发送消息的过程中可能会出现异常,或者网络的抖动,故障等等因为造成消息的无法消费,比如用户在下订单,消费MQ接受,订单系统出现故障,导致用户支付失败,那么这个时候就需要消息中间件就必须支持消息重试机制策略。也就是支持:出现问题和故障的情况下,消息不丢失还可以进行重发。

6.4、消费分发策略的机制和对比

  ActiveMQ RabbitMQ kafka RocketMQ
发布订阅 支持 支持 支持 支持
轮询分发 支持 支持 支持 /
公平分发 / 支持 支持 /
重发 支持 支持 / 支持
消息拉取 / 支持 支持 支持

7、消息队列的高可用和高可靠

7.1、什么是高可用机制

所谓高可用:是指产品在规定的条件和规定的时刻或时间内处于可执行规定功能状态的能力。 当业务量增加时,请求也过大,一台消息中间件服务器的会触及硬件(CPU,内存,磁盘)的极限,一台消息服务器你已经无法满足业务的需求,所以消息中间件必须支持集群部署。来达到高可用的目的。

7.2、集群模式1:master-slave主从共享数据的部署方式

image-20230515230107475

解说:生产者将消息发送到Master节点,所有的都连接这个消息队列共享这块数据区域,Master节点负责写入,一旦Master挂掉,slave节点继续服务。从而形成高可用,

7.3、集群模式2:master-slave主从同步部署方式

image-20230515230210308

解释:这种模式写入消息同样在Master主节点上,但是主节点会同步数据到slave节点形成副本,和zookeeper或者redis主从机制很类同。这样可以达到负载均衡的效果,如果消费者有多个这样就可以去不同的节点就行消费,以为消息的拷贝和同步会占用很大的带宽和网络资源。在后续的rabbtmq中会有使用。

7.4、集群模式3:多主集群同步部署方式

image-20230515230358993

解释:和上面的区别不是特别的大,但是它的写入可以往任意节点去写入。

7.5、集群模式4:多主集群转发部署方式

image-20230515230416331

解释:如果你插入的数据是broker-1中,元数据信息会存储数据的相关描述和记录存放的位置(队列)。 它会对描述信息也就是元数据信息就行同步,如果消费者在broker-2中进行消费,发现自己几点没有对应的消息,可以从对应的元数据信息中去查询,然后返回对应的消息信息,场景:比如买火车票或者黄牛买演唱会门票,比如第一个黄牛有顾客说要买的演唱会门票,但是没有但是他会去联系其他的黄牛询问,如果有就返回。

7.6、集群模式5:master-slave与breoker-cluster组合的方案

image-20230515230701227

解释:实现多主多从的热备机制来完成消息的高可用以及数据的热备机制,在生产规模达到一定的阶段的时候,这种使用的频率比较高。

这么集群模式,具体在后续的课程中会进行一个分析和讲解。他们的最终目的都是为保证:消息服务器不会挂掉,出现了故障依然可以抱着消息服务继续使用。

反正终归三句话:

  • 要么消息共享
  • 要么消息同步
  • 要么元数据共享

7.7、什么是高可靠机制

指系统可以无故障低持续运行,比如一个系统突然崩溃,报错,异常等等并不影响线上业务的正常运行,出错的几率极低,就称之为:高可靠

在高并发的业务场景中,如果不能保证系统的高可靠,那造成的隐患和损失是非常严重的。

如何保证中间件消息的可靠性呢?可以从两个方面考虑

  1. 消息的传输:通过协议来保证系统间数据解析的正确性
  2. 消息的存储可靠:通过持久化来保证消息的可靠性

8、rabbitMQ入门及安装

8.1、概述

官网:https://www.rabbitmq.com/ 什么是RabbitMQ,官方给出来这样的解释:

RabbitMQ is the most widely deployed open source message broker.
With tens of thousands of users, RabbitMQ is one of the most popular open source message brokers. From T-Mobile to Runtastic, RabbitMQ is used worldwide at small startups and large enterprises.
RabbitMQ is lightweight and easy to deploy on premises and in the cloud. It supports multiple messaging protocols. RabbitMQ can be deployed in distributed and federated configurations to meet high-scale, high-availability requirements.
RabbitMQ runs on many operating systems and cloud environments, and provides a wide range of developer tools for most popular languages.
翻译以后:
RabbitMQ是部署最广泛的开源消息代理。
RabbitMQ拥有成千上万的用户,是最受欢迎的开源消息代理之一。从T-Mobile 到Runtastic,RabbitMQ在全球范围内的小型初创企业和大型企业中都得到使用。
RabbitMQ轻巧,易于在内部和云中部署。它支持多种消息传递协议。RabbitMQ可以部署在分布式和联合配置中,以满足大规模,高可用性的要求。
RabbitMQ可在许多操作系统和云环境上运行,并为大多数流行语言提供了广泛的开发人员工具。

简单概述: RabbitMQ是一个开源的遵循AMQP协议实现的基于Erlang语言编写,支持多种客户端(语言)。用于在分布式系统中存储消息,转发消息,具有高可用,高可扩性,易用性等特征。

8.2、安装rabbitMQ

  1. 查看系统版本号

    [root@iZwz9efdd2ukk4oauustczZ rabbitmq]# lsb_release -a
    LSB Version:	:core-4.1-amd64:core-4.1-noarch
    Distributor ID:	CentOS
    Description:	CentOS Linux release 7.6.1810 (Core) 
    Release:	7.6.1810
    Codename:	Core
    
  2. 先下载Erlang

    下载地址:https://github.com/rabbitmq/erlang-rpm/releases/download/v23.2.6/erlang-23.2.6-1.el7.x86_64.rpm

  3. 上传到服务器上

    image-20230517172454500

  4. 安装并查看

    [root@iZwz9efdd2ukk4oauustczZ rabbitmq]# rpm -Uvh erlang-23.2.6-1.el7.x86_64.rpm  # 安装
    warning: erlang-23.2.6-1.el7.x86_64.rpm: Header V4 RSA/SHA1 Signature, key ID 6026dfca: NOKEY
    Preparing...                          ################################# [100%]
    Updating / installing...
       1:erlang-23.2.6-1.el7              ################################# [100%]
    [root@iZwz9efdd2ukk4oauustczZ rabbitmq]# erl -version # 查看版本,Erlang安装成功
    Erlang (SMP,ASYNC_THREADS,HIPE) (BEAM) emulator version 11.1.8
    
  5. 下载rabbitMQ,注意与erlang的版本对应

    下载地址:https://github.com/rabbitmq/rabbitmq-server/releases/download/v3.8.12/rabbitmq-server-3.8.12-1.el7.noarch.rpm

    erlang和RabbitMQ版本的按照比较: https://www.rabbitmq.com/which-erlang.html

    image-20230517172811544

  6. 上传到服务器上

  7. 安装

    [root@iZwz9efdd2ukk4oauustczZ rabbitmq]# yum install -y rabbitmq-server-3.8.12-1.el7.noarch.rpm # 安装
    Loaded plugins: fastestmirror
    Examining rabbitmq-server-3.8.12-1.el7.noarch.rpm: rabbitmq-server-3.8.12-1.el7.noarch
    Marking rabbitmq-server-3.8.12-1.el7.noarch.rpm to be installed
    Resolving Dependencies
    --> Running transaction check
    ---> Package rabbitmq-server.noarch 0:3.8.12-1.el7 will be installed
    --> Processing Dependency: socat for package: rabbitmq-server-3.8.12-1.el7.noarch
    Determining fastest mirrors
    base                                                                                                                                                   | 3.6 kB  00:00:00     
    docker-ce-stable                                                                                                                                       | 3.5 kB  00:00:00     
    epel                                                                                                                                                   | 4.7 kB  00:00:00     
    extras                                                                                                                                                 | 2.9 kB  00:00:00     
    updates                                                                                                                                                | 2.9 kB  00:00:00     
    (1/4): epel/x86_64/updateinfo                                                                                                                          | 1.0 MB  00:00:00     
    (2/4): epel/x86_64/primary_db                                                                                                                          | 7.0 MB  00:00:00     
    (3/4): updates/7/x86_64/primary_db                                                                                                                     |  21 MB  00:00:00     
    (4/4): docker-ce-stable/7/x86_64/primary_db                                                                                                            | 109 kB  00:00:00     
    --> Running transaction check
    ---> Package socat.x86_64 0:1.7.3.2-2.el7 will be installed
    --> Finished Dependency Resolution
    
    Dependencies Resolved
    
    ==============================================================================================================================================================================
     Package                                Arch                          Version                               Repository                                                   Size
    ==============================================================================================================================================================================
    Installing:
     rabbitmq-server                        noarch                        3.8.12-1.el7                          /rabbitmq-server-3.8.12-1.el7.noarch                         16 M
    Installing for dependencies:
     socat                                  x86_64                        1.7.3.2-2.el7                         base                                                        290 k
    
    Transaction Summary
    ==============================================================================================================================================================================
    Install  1 Package (+1 Dependent package)
    
    Total size: 16 M
    Total download size: 290 k
    Installed size: 17 M
    Downloading packages:
    socat-1.7.3.2-2.el7.x86_64.rpm                                                                                                                         | 290 kB  00:00:00     
    Running transaction check
    Running transaction test
    Transaction test succeeded
    Running transaction
    Warning: RPMDB altered outside of yum.
      Installing : socat-1.7.3.2-2.el7.x86_64                                                                                                                                 1/2 
      Installing : rabbitmq-server-3.8.12-1.el7.noarch                                                                                                                        2/2 
      Verifying  : socat-1.7.3.2-2.el7.x86_64                                                                                                                                 1/2 
      Verifying  : rabbitmq-server-3.8.12-1.el7.noarch                                                                                                                        2/2 
    
    Installed:
      rabbitmq-server.noarch 0:3.8.12-1.el7                                                                                                                                       
    
    Dependency Installed:
      socat.x86_64 0:1.7.3.2-2.el7                                                                                                                                                
    
    Complete!
    

8.3、启动、停止rabbitMQ服务

# 启动服务
> systemctl start rabbitmq-server
# 查看服务状态
> systemctl status rabbitmq-server
# 停止服务
> systemctl stop rabbitmq-server
# 开机启动服务
> systemctl enable rabbitmq-server

8.4、rabbitMQ配置

RabbitMQ默认情况下有一个配置文件,定义了RabbitMQ的相关配置信息,默认情况下能够满足日常的开发需求。如果需要修改需要,需要自己创建一个配置文件进行覆盖。

参考官网:

  1. https://www.rabbitmq.com/documentation.html
  2. https://www.rabbitmq.com/configure.html
  3. https://www.rabbitmq.com/configure.html#config-items
  4. https://github.com/rabbitmq/rabbitmq-server/blob/add-debug-messages-to-quorum_queue_SUITE/docs/rabbitmq.conf.example

8.5、相关端口

  • 5672:RabbitMQ的通讯端口
  • 25672:RabbitMQ的节点间的CLI通讯端口
  • 15672:RabbitMQ HTTP_API的端口,管理员用户才能访问,用于管理RabbitMQ,需要启动Management插件
  • 1883,8883:MQTT插件启动时的端口
  • 61613、61614:STOMP客户端插件启用的时候的端口
  • 15674、15675:基于webscoket的STOMP端口和MOTT端口

一定要注意:RabbitMQ 在安装完毕以后,会绑定一些端口,如果你购买的是阿里云或者腾讯云相关的服务器一定要在安全组中把对应的端口添加到防火墙。

9、rabbitMQ管理页面及授权操作

9.1、rabbitMQ管理页面

  1. 默认情况下,rabbitmq是没有安装web端的客户端插件,需要安装才可以生效

    [root@iZwz9efdd2ukk4oauustczZ rabbitmq]# rabbitmq-plugins enable rabbitmq_management
    Enabling plugins on node rabbit@iZwz9efdd2ukk4oauustczZ:
    rabbitmq_management
    The following plugins have been configured:
      rabbitmq_management
      rabbitmq_management_agent
      rabbitmq_web_dispatch
    Applying plugin configuration to rabbit@iZwz9efdd2ukk4oauustczZ...
    The following plugins have been enabled:
      rabbitmq_management
      rabbitmq_management_agent
      rabbitmq_web_dispatch
    
    started 3 plugins.
    

    说明:rabbitmq有一个默认账号和密码是:guest 默认情况只能在localhost本机下访问,所以需要添加一个远程登录的用户。

  2. 安装完毕以后,重启服务即可

    [root@iZwz9efdd2ukk4oauustczZ rabbitmq]# systemctl restart rabbitmq-server
    

    一定要记住,在对应服务器(阿里云,腾讯云等)的安全组中开放15672的端口。

    image-20230517180110440

  3. 在浏览器访问

    http://39.108.49.252:15672

    image-20230517180031469

9.2、授权账户和密码

rabbitmqctl add_user 账号 密码  # 添加账号
rabbitmqctl set_user_tags 账号 administrator # 设置账号的级别
rabbitmqctl change_password Username Newpassword # 修改密码 
rabbitmqctl delete_user Username # 删除用户
rabbitmqctl list_users # 查看用户清单
rabbitmqctl set_permissions -p / 用户名 ".*" ".*" ".*" # 为用户设置administrator角色

用户级别

  • administrator 可以登录控制台、查看所有信息、可以对rabbitmq进行管理
  • monitoring 监控者 登录控制台,查看所有信息
  • policymaker 策略制定者 登录控制台,指定策略
  • managment 普通管理员 登录控制台

测试

  1. 新增用户

    [root@iZwz9efdd2ukk4oauustczZ rabbitmq]# rabbitmqctl add_user admin admin
    Adding user "admin" ...
    Done. Don't forget to grant the user permissions to some virtual hosts! See 'rabbitmqctl help set_permissions' to learn more.
    
  2. 设置用户分配操作权限

    [root@iZwz9efdd2ukk4oauustczZ rabbitmq]# rabbitmqctl set_user_tags admin administrator
    Setting tags for user "admin" to [administrator] ...
    
  3. 为用户添加资源权限

    [root@iZwz9efdd2ukk4oauustczZ rabbitmq]# rabbitmqctl set_permissions -p / admin ".*" ".*" ".*"
    Setting permissions for user "admin" in vhost "/" ...
    

然后就可以登录管理页面了

image-20230517181045742

10、rabbitMQ之docker安装

10.1、docker安装

已经安装docker的,可以忽略这一步

(1)yum 包更新到最新
> yum update
(2)安装需要的软件包, yum-util 提供yum-config-manager功能,另外两个是devicemapper驱动依赖的
> yum install -y yum-utils device-mapper-persistent-data lvm2
(3)设置yum源为阿里云
> yum-config-manager --add-repo http://mirrors.aliyun.com/docker-ce/linux/centos/docker-ce.repo
(4)安装docker
> yum install docker-ce -y
(5)安装后查看docker版本
> docker -v
 (6) 安装加速镜像
 sudo mkdir -p /etc/docker
 sudo tee /etc/docker/daemon.json <<-'EOF'
 {
  "registry-mirrors": ["https://0wrdwnn6.mirror.aliyuncs.com"]
 }
 EOF
 sudo systemctl daemon-reload
 sudo systemctl restart docker

docker的相关命令

# 启动docker:
systemctl start docker
# 停止docker:
systemctl stop docker
# 重启docker:
systemctl restart docker
# 查看docker状态:
systemctl status docker
# 开机启动:  
systemctl enable docker
systemctl unenable docker
# 查看docker概要信息
docker info
# 查看docker帮助文档
docker --help

10.2、安装rabbitMQ

  1. 获取rabbitMQ镜像

    [root@iZwz9efdd2ukk4oauustczZ ~]# docker pull rabbitmq:management
    management: Pulling from library/rabbitmq
    7b1a6ab2e44d: Pull complete 
    37f453d83d8f: Pull complete 
    e64e769bc4fd: Pull complete 
    c288a913222f: Pull complete 
    12addf9c8bf9: Pull complete 
    eaeb088e057d: Pull complete 
    b63d48599313: Pull complete 
    05c99d3d2a57: Pull complete 
    43665bfbc3f9: Pull complete 
    f14c7d7911b1: Pull complete 
    Digest: sha256:4c4b66ad5ec40b2c27943b9804d307bf31c17c8537cd0cd107236200a9cd2814
    Status: Downloaded newer image for rabbitmq:management
    docker.io/library/rabbitmq:management
    

    image-20230517223602431

  2. 创建并运行容器

    docker run -di --name=myrabbit -p 15672:15672 rabbitmq:management
    

    运行时设置用户和密码,命令如下

    [root@iZwz9efdd2ukk4oauustczZ ~]# docker run -di --name myrabbitmq -e RABBITMQ_DEFAULT_USER=admin -e RABBITMQ_DEFAULT_PASS=admin -p 15672:15672 -p 5672:5672 -p 25672:25672 -p 61613:61613 -p 1883:1883 rabbitmq:management
    f1d06bf3a6015580e228df8bb54348a0b82d5ee3c76959c306829a2fad3c46fe
    
  3. 查看日志

    [root@iZwz9efdd2ukk4oauustczZ ~]# docker logs -f myrabbitmq
    
  4. 容器运行正常

    请求http://39.108.49.252:15672

    image-20230517223757730

拓展linux命令:

> more xxx.log  #查看日记信息
> netstat -naop | grep 5672 #查看端口是否被占用
> ps -ef | grep 5672  #查看进程
> systemctl stop 服务 # 停止服务

11、rabbitMQ的角色分类

  1. none
    • 不能访问management plugin
  2. management 查看自己相关节点信息
    • 列出自己可以通过AMQP登入的虚拟机
    • 查看自己的虚拟机节点 virtual hosts的queues,exchanges和bindings信息
    • 查看和关闭自己的channels和connections
    • 查看有关自己的虚拟机节点virtual hosts的统计信息。包括其他用户在这个节点virtual hosts中的活动信息。
  3. policymaker
    • 包含management所有权限
    • 查看和创建和删除自己的virtual hosts所属的policies和parameters信息
  4. monitoring
    • 包含management所有权限
    • 罗列出所有的virtual hosts,包括不能登录的virtual hosts
    • 查看其他用户的connections和channels信息
    • 查看节点级别的数据如clustering和memory使用情况
    • 查看所有的virtual hosts的全局统计信息
  5. administrator
    • 最高权限
    • 可以创建和删除virtual hosts
    • 可以查看,创建和删除users
    • 查看创建permisssions
    • 关闭所有用户的connections

12、rabbitMQ基本概念

12.1、生产者、消费者

  • 生产者(Producer)

    消息的创建者

    负责创建和推送数据到消息服务器

  • 消费者(Consumer)

    消息的接收方

    负责接收消息和处理数据

12.2、消息队列

消息队列是RabbitMQ的内部对象,用于存储生产者的消息直到发送给消费者,它是消费者接收消息的地方。

12.3、交换机

交换机用于接收,分配消息

交换机包含4中类型: direct, topic, fanout, headers。

image-20230606234613515

  • direct(直连交换机)

    具有路由功能的交换机,绑定到此交换机的时候需要指定一个routing_key,交换机发送消息的时候需要routing_key,会将消息发送道对应的队列

    先匹配,再投送

    Direct Exchange是RabbitMQ的默认交换机模式

    这是最简单的模式

    它根据routing key全文匹配去寻找队列

  • topic(主题交换机)

    在直连交换机基础上增加模式匹配,也就是对routing_key进行模式匹配,*代表一个单词,#代表多个单词

    按规则转发消息

    主题交换机(Topic Exchange)主要根据通配符转发消息

    这种方式最灵活

    交换机和队列的绑定会定义一种路由模式

    路由键(routing key)和路由模式匹配后,交换机才能转发消息

    在这种交换机模式下,路由键(routing key)必须是一串字符,用"."隔开
    路由模式必须包含一个星号"*", 主要用于匹配路由键指定位置的一个单词
    * 匹配一个单词
    # 匹配0个或多个单词
    
    eg:
    binding key:                 *.com.#
    匹配的routing key:     cn.com,  us.com.aa
    不匹配:                         com.bb
    
  • fanout(扇形交换机)

    广播消息到所有队列,没有任何处理,速度最快

    消息广播的模式

    这种方式将消息广播到所有绑定到它的队列中。 不考虑routing key的值,即使配置了路由键,依然会被忽略。

  • headers(首部交换机)

    忽略routing_key,使用Headers信息(一个Hash的数据结构)进行匹配,优势在于可以有更多更灵活的匹配规则

    根据应用程序消息的特定属性进行匹配

13、rabbitMQ七种工作模式

参考官网:https://www.rabbitmq.com/getstarted.html

13.1、简单模式(Hello World)

image-20230606233242623

做最简单的事情,一个生产者对应一个消费者,RabbitMQ相当于一个消息代理,负责将A的消息转发给B

单生产者,单消费者,单队列

13.2、工作队列模式(Work queues)

image-20230606233422844

在多个消费者之间分配任务(竞争的消费者模式),一个生产者对应多个消费者。

适用于资源密集型任务, 单个消费者处理不过来,需要多个消费者进行处理的场景。

单生产者,多消费者,单队列。

应用场景:

一个订单的处理需要10s,有多个订单可以同时放到消息队列,

然后让多个消费者同时并行处理,而不是单个消费者的串行消费。

13.3、发布订阅模式(Publish/Subscribe)

image-20230606233545406

一次向许多消费者发送消息,将消息将广播到所有的消费者。

单生产者,多消费者,多队列

应用场景:

更新商品库存后需要通知多个缓存和多个数据库。

结构如下:

  • 一个fanout类型交换机扇出两个消息队列,分别为缓存消息队列、数据库消息队列
  • 一个缓存消息队列对应着多个缓存消费者
  • 一个数据库消息队列对应着多个数据库消费者

13.4、路由模式(Routing)

image-20230606233810239

根据Routing Key有选择地接收消息。

多消费者,选择性多队列,每个队列通过routing key全文匹配。

发送消息到交换机并且要指定路由键(Routing key) 。 消费者将队列绑定到交换机时需要指定路由key,仅消费指定路由key的消息

应用场景:

在商品库存中增加了1台iphone12,iphone12促销活动消费者指定routing key为iphone12, 只有此促销活动会接收到消息,其它促销活动不关心也不会消费此routing key的消息。

13.5、主题模式(Topics)

image-20230606234010389

主题交换机方式接收消息,将routing key和模式进行匹配。

多消费者,选择性多队列,每个队列通过模式匹配。

队列需要绑定在一个模式上。 #匹配一个词或多个词,*只匹配一个词。

应用场景:

iphone促销活动可以接收主题为多种iPhone的消息,如iphone12、iphone13等。

13.6、远程过程调用(RPC)

image-20230606234108384

在远程计算机上运行功能并等待结果。

应用场景:

需要等待接口返回数据,如订单支付。

13.7、发布者确认(Publisher Confirms)

与发布者进行可靠的发布确认,发布者确认是RabbitMQ扩展,可以实现可靠的发布。

在通道上启用发布者确认后,RabbitMQ将异步确认发送者发布的消息,这意味着它们已在服务器端处理。

应用场景:

对于消息可靠性要求较高,比如钱包扣款。

14、rabbitMQ入门案例-simple简单模式

  1. 构建一个maven项目rabbitmq-java

  2. 导包

    <?xml version="1.0" encoding="UTF-8"?>
    <project xmlns="http://maven.apache.org/POM/4.0.0"
             xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
             xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
        <modelVersion>4.0.0</modelVersion>
    
        <groupId>com.zyy</groupId>
        <artifactId>rabbitmq-java</artifactId>
        <version>1.0-SNAPSHOT</version>
    
        <dependencies>
            <dependency>
                <groupId>com.rabbitmq</groupId>
                <artifactId>amqp-client</artifactId>
                <version>5.10.0</version>
            </dependency>
        </dependencies>
    </project>
    
  3. 创建好目录,并新建类Producer

    import com.rabbitmq.client.Channel;
    import com.rabbitmq.client.Connection;
    import com.rabbitmq.client.ConnectionFactory;
    
    import java.io.IOException;
    import java.util.concurrent.TimeoutException;
    
    public class Producer {
        public static void main(String[] args) {
            //1.创建连接工厂
            ConnectionFactory connectionFactory = new ConnectionFactory();
            //2.设置连接属性
            connectionFactory.setHost("39.108.49.252");
            connectionFactory.setPort(5672);
            connectionFactory.setUsername("admin");
            connectionFactory.setPassword("admin");
            connectionFactory.setVirtualHost("/");
            Connection connection = null;
            Channel channel = null;
            try {
                //3.从连接工厂中获取连接
                connection = connectionFactory.newConnection("zyy_producer");
                //4.从连接中创建通道channel
                channel = connection.createChannel();
                /**
                 * 5.声明队列queue
                 * @param queue 队列的名称
                 * @param durable 队列是否持久化
                 * @param exclusive 是否排他,即是否私有化,如果为true,会对当前队列加锁,其他的通道不能访问,并且连接自动关闭
                 * @param autoDelete 是否自动删除,当最后一个消费者断开连接之后是否自动删除消息
                 * @param arguments 可以设置队列的附件参数,设置队列的有效期,消息的最大长度,队列的消息生命周期等
                 */
                channel.queueDeclare("zyy_queue", false, false, false, null);
                //6.发送消息
                String message = "zyy hello!";
                /**
                 *
                 * @param exchange 交换机exchange
                 * @param routingKey 队列对称
                 * @param props 属性配置
                 * @param body 发送消息的内容
                 */
                channel.basicPublish("", "zyy_queue", null, message.getBytes());
                System.out.println("消息发送成功");
            } catch (IOException e) {
                e.printStackTrace();
            } catch (TimeoutException e) {
                e.printStackTrace();
            } finally {
                if (channel != null && channel.isOpen()) {
                    try {
                        channel.close();
                    } catch (IOException e) {
                        e.printStackTrace();
                    } catch (TimeoutException e) {
                        e.printStackTrace();
                    }
                }
                //7.释放连接
                if (connection != null) {
                    try {
                        connection.close();
                    } catch (IOException e) {
                        e.printStackTrace();
                    }
                }
            }
        }
    }
    
    

    image-20230518232949680

  4. 运行(前提需要启动服务器上的rabbitMQ服务),查看管理页面

    image-20230518232637203

    image-20230518232713560

    页面显示描述

    image-20230518233606781

    image-20230518233649212

  5. 继续定义消费者

    import com.rabbitmq.client.CancelCallback;
    import com.rabbitmq.client.Channel;
    import com.rabbitmq.client.Connection;
    import com.rabbitmq.client.ConnectionFactory;
    import com.rabbitmq.client.DeliverCallback;
    import com.rabbitmq.client.Delivery;
    
    import java.io.IOException;
    import java.util.concurrent.TimeoutException;
    
    public class Consumer {
        public static void main(String[] args) {
            //1.创建连接工厂
            ConnectionFactory connectionFactory = new ConnectionFactory();
            //2.设置连接属性
            connectionFactory.setHost("39.108.49.252");
            connectionFactory.setPort(5672);
            connectionFactory.setUsername("admin");
            connectionFactory.setPassword("admin");
            connectionFactory.setVirtualHost("/");
            Connection connection = null;
            Channel channel = null;
            try {
                //3.从连接工厂中获取连接
                connection = connectionFactory.newConnection("zyy_producer");
                //4.从连接中创建通道channel
                channel = connection.createChannel();
                /**
                 * 5.消费信息
                 */
                channel.basicConsume("zyy_queue", true, new DeliverCallback() {
                    public void handle(String consumerTag, Delivery message) throws IOException {
                        System.out.println(consumerTag + "收到的消息是" + new String(message.getBody(), "UTF-8"));
                    }
                }, new CancelCallback() {
                    public void handle(String consumerTag) throws IOException {
                        System.out.println(consumerTag + "接受失败了");
                    }
                });
    
            } catch (IOException e) {
                e.printStackTrace();
            } catch (TimeoutException e) {
                e.printStackTrace();
            } finally {
                //6.释放连接
                if (channel != null && channel.isOpen()) {
                    try {
                        channel.close();
                    } catch (IOException e) {
                        e.printStackTrace();
                    } catch (TimeoutException e) {
                        e.printStackTrace();
                    }
                }
                if (connection != null) {
                    try {
                        connection.close();
                    } catch (IOException e) {
                        e.printStackTrace();
                    }
                }
            }
        }
    }
    
    
  6. 观察管理页面,发现消息被消费了

    image-20230518234559851

15、什么是AMQP

什么是AMQP

AMQP全称:Advanced Message Queuing Protocol(高级消息队列协议)。是应用层协议的一个开发标准,为面向消息的中间件设计。

AMQP生产者流转过程

image-20230521163611162

AMQP消费者流转过程

image-20230521163837479

16、rabbitMQ的核心组成部分

16.1、rabbitMQ的核心组成部分

image-20230521164322592

  • Server:又称Broker ,接受客户端的连接,实现AMQP实体服务。 安装rabbitmq-server
  • Connection:连接,应用程序与Broker的网络连接 TCP/IP 三次握手和四次挥手
  • Channel:网络信道,几乎所有的操作都在Channel中进行,Channel是进行消息读写的通道,客户端可以建立对各Channel,每个Channel代表一个会话任务
  • Message :消息:服务与应用程序之间传送的数据,由Properties和body组成,Properties可是对消息进行修饰,比如消息的优先级,延迟等高级特性,Body则就是消息体的内容
  • Virtual Host 虚拟地址,用于进行逻辑隔离,最上层的消息路由,一个虚拟主机可以有若干个Exhange和Queueu,同一个虚拟主机里面不能有相同名字的Exchange
  • Exchange:交换机,接受消息,根据路由键发送消息到绑定的队列。(不具备消息存储的能力)
  • Bindings:Exchange和Queue之间的虚拟连接,binding中可以保护多个routing key
  • Routing key:是一个路由规则,虚拟机可以用它来确定如何路由一个特定消息
  • Queue:队列:也成为Message Queue,消息队列,保存消息并将它们转发给消费者

16.2、rabbitMQ整体架构是什么样子的

image-20230521164810513

16.3、rabbitMQ的运行流程

image-20230521164909723

16.4、小结

  • rabbitmq发送消息一定有一个交换机
  • image-20230521165214512
  • image-20230521165228471

17、rabbitMQ入门案例-fanout交换机-发布订阅模式

image-20230521233257064

  1. 新建四个队列

    image-20230521234239438

  2. 新建一个fanout类型的交换机

    image-20230521234322147

  3. 绑定queue1、queue2

    image-20230521234401500

  4. 发送消息

    image-20230521234458244

  5. queue1、queue2都收到了

    image-20230521234550915

java代码实现

  1. 生产者

    import com.rabbitmq.client.Channel;
    import com.rabbitmq.client.Connection;
    import com.rabbitmq.client.ConnectionFactory;
    
    import java.io.IOException;
    import java.util.concurrent.TimeoutException;
    
    public class Producer {
    
        public static void main(String[] args) {
            //创建连接并设置连接属性
            ConnectionFactory connectionFactory = new ConnectionFactory();
            connectionFactory.setVirtualHost("/");
            connectionFactory.setHost("39.108.49.252");
            connectionFactory.setPort(5672);
            connectionFactory.setUsername("admin");
            connectionFactory.setPassword("admin");
    
            Connection connection = null;
            Channel channel = null;
            try {
                //从连接工厂中新建连接
                connection = connectionFactory.newConnection("producer-1");
                //从连接中创建通道
                channel = connection.createChannel();
                //发送消息
                String exchangeName = "fanout.exchange";
                String message = "hello fanout.exchange";
                channel.basicPublish(exchangeName, "", null, message.getBytes());
                System.out.println("消息发送成功");
            } catch (IOException e) {
                e.printStackTrace();
            } catch (TimeoutException e) {
                e.printStackTrace();
            } finally {
                //关闭通道和连接
                if (channel != null && channel.isOpen()) {
                    try {
                        channel.close();
                    } catch (IOException e) {
                        e.printStackTrace();
                    } catch (TimeoutException e) {
                        e.printStackTrace();
                    }
                }
    
                if (connection != null && connection.isOpen()) {
                    try {
                        connection.close();
                    } catch (IOException e) {
                        e.printStackTrace();
                    }
                }
            }
        }
    }
    
  2. 消费者

    import com.rabbitmq.client.CancelCallback;
    import com.rabbitmq.client.Channel;
    import com.rabbitmq.client.Connection;
    import com.rabbitmq.client.ConnectionFactory;
    import com.rabbitmq.client.DeliverCallback;
    import com.rabbitmq.client.Delivery;
    
    import java.io.IOException;
    import java.util.concurrent.TimeoutException;
    
    
    public class Consumer {
    
        private static Runnable runnable = () -> {
            //创建连接并设置连接属性
            ConnectionFactory connectionFactory = new ConnectionFactory();
            connectionFactory.setVirtualHost("/");
            connectionFactory.setHost("39.108.49.252");
            connectionFactory.setPort(5672);
            connectionFactory.setUsername("admin");
            connectionFactory.setPassword("admin");
    
            String queueName = Thread.currentThread().getName();
    
            Connection connection = null;
            Channel channel = null;
            try {
                //从连接工厂中新建连接
                connection = connectionFactory.newConnection("producer-1");
                //从连接中创建通道
                channel = connection.createChannel();
    
                //接受消息
                channel.basicConsume(queueName, true, new DeliverCallback() {
                    @Override
                    public void handle(String consumerTag, Delivery message) throws IOException {
                        System.out.println(queueName + ",收到信息是:" + new String(message.getBody(), "UTF-8"));
                    }
                }, new CancelCallback() {
                    @Override
                    public void handle(String consumerTag) throws IOException {
    
                    }
                });
                System.out.println(queueName + ",开始接受消息");
                System.in.read();
            } catch (IOException e) {
                e.printStackTrace();
            } catch (TimeoutException e) {
                e.printStackTrace();
            } finally {
                //关闭通道和连接
                if (channel != null && channel.isOpen()) {
                    try {
                        channel.close();
                    } catch (IOException e) {
                        e.printStackTrace();
                    } catch (TimeoutException e) {
                        e.printStackTrace();
                    }
                }
    
                if (connection != null && connection.isOpen()) {
                    try {
                        connection.close();
                    } catch (IOException e) {
                        e.printStackTrace();
                    }
                }
            }
        };
    
        public static void main(String[] args) {
            new Thread(runnable, "queue1").start();
            new Thread(runnable, "queue2").start();
        }
    
    }
    
  3. 运行生产者,发送消息,运行消费者,消费消息

    queue1,收到信息是:hello fanout.exchange
    queue2,收到信息是:hello fanout.exchange
    

上面的例子是在管理页面做的:新建交换机、新建队列,绑定队列;java代码做的:生产消息和消费消息,那我们试试都在代码中实现

只有生产者不一样,消费者一样的(ps:改一下main方法,添加queue5,queue6子线程即可),生产者代码如下:

import com.rabbitmq.client.BuiltinExchangeType;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;

import java.io.IOException;
import java.util.concurrent.TimeoutException;

public class Producer {

    public static void main(String[] args) {
        //创建连接并设置连接属性
        ConnectionFactory connectionFactory = new ConnectionFactory();
        connectionFactory.setVirtualHost("/");
        connectionFactory.setHost("39.108.49.252");
        connectionFactory.setPort(5672);
        connectionFactory.setUsername("admin");
        connectionFactory.setPassword("admin");

        Connection connection = null;
        Channel channel = null;
        try {
            //从连接工厂中新建连接
            connection = connectionFactory.newConnection("producer-1");
            //从连接中创建通道
            channel = connection.createChannel();
            //声明交换机
            String exchangeName = "fanout.java.exchange";
            channel.exchangeDeclare(exchangeName, BuiltinExchangeType.FANOUT, true);
            //声明队列
            String queueName1 = "queue5";
            String queueName2 = "queue6";
            channel.queueDeclare(queueName1, true, false, false, null);
            channel.queueDeclare(queueName2, true, false, false, null);
            //绑定队列
            channel.queueBind(queueName1, exchangeName, "", null);
            channel.queueBind(queueName2, exchangeName, "", null);
            //发送消息
            String message = "hello fanout.java.exchange";
            channel.basicPublish(exchangeName, "", null, message.getBytes());
            System.out.println("消息发送成功");
        } catch (IOException e) {
            e.printStackTrace();
        } catch (TimeoutException e) {
            e.printStackTrace();
        } finally {
            //关闭通道和连接
            if (channel != null && channel.isOpen()) {
                try {
                    channel.close();
                } catch (IOException e) {
                    e.printStackTrace();
                } catch (TimeoutException e) {
                    e.printStackTrace();
                }
            }

            if (connection != null && connection.isOpen()) {
                try {
                    connection.close();
                } catch (IOException e) {
                    e.printStackTrace();
                }
            }
        }
    }
}

18、rabbitMQ入门案例-direct交换机-路由模式

image-20230521233310868

  1. 在新建一个direct交换机

    image-20230521234742228

  2. 绑定交换机,并指定routing key

    image-20230521235002680

  3. 发送消息,指定routing key = key1

    image-20230521235103003

  4. 发现只有指定的队列可以收到

    image-20230521235156450

java代码实现

  1. 生产者

    import com.rabbitmq.client.Channel;
    import com.rabbitmq.client.Connection;
    import com.rabbitmq.client.ConnectionFactory;
    
    import java.io.IOException;
    import java.util.concurrent.TimeoutException;
    
    public class Producer {
    
        public static void main(String[] args) {
            //创建连接并设置连接属性
            ConnectionFactory connectionFactory = new ConnectionFactory();
            connectionFactory.setVirtualHost("/");
            connectionFactory.setHost("39.108.49.252");
            connectionFactory.setPort(5672);
            connectionFactory.setUsername("admin");
            connectionFactory.setPassword("admin");
    
            Connection connection = null;
            Channel channel = null;
            try {
                //从连接工厂中新建连接
                connection = connectionFactory.newConnection("producer-1");
                //从连接中创建通道
                channel = connection.createChannel();
                //发送消息
                String exchangeName = "direct.exchange";
                String message = "hello direct.exchange";
                //指定了routingKey
                String routingKey = "key1";
                channel.basicPublish(exchangeName, routingKey, null, message.getBytes());
                System.out.println("消息发送成功");
            } catch (IOException e) {
                e.printStackTrace();
            } catch (TimeoutException e) {
                e.printStackTrace();
            } finally {
                //关闭通道和连接
                if (channel != null && channel.isOpen()) {
                    try {
                        channel.close();
                    } catch (IOException e) {
                        e.printStackTrace();
                    } catch (TimeoutException e) {
                        e.printStackTrace();
                    }
                }
    
                if (connection != null && connection.isOpen()) {
                    try {
                        connection.close();
                    } catch (IOException e) {
                        e.printStackTrace();
                    }
                }
            }
        }
    }
    
    
  2. 消费者

    import com.rabbitmq.client.CancelCallback;
    import com.rabbitmq.client.Channel;
    import com.rabbitmq.client.Connection;
    import com.rabbitmq.client.ConnectionFactory;
    import com.rabbitmq.client.DeliverCallback;
    import com.rabbitmq.client.Delivery;
    
    import java.io.IOException;
    import java.util.concurrent.TimeoutException;
    
    public class Consumer {
    
        private static Runnable runnable = () -> {
            //创建连接并设置连接属性
            ConnectionFactory connectionFactory = new ConnectionFactory();
            connectionFactory.setVirtualHost("/");
            connectionFactory.setHost("39.108.49.252");
            connectionFactory.setPort(5672);
            connectionFactory.setUsername("admin");
            connectionFactory.setPassword("admin");
    
            String queueName = Thread.currentThread().getName();
    
            Connection connection = null;
            Channel channel = null;
            try {
                //从连接工厂中新建连接
                connection = connectionFactory.newConnection("producer-1");
                //从连接中创建通道
                channel = connection.createChannel();
    
                //接受消息
                channel.basicConsume(queueName, true, new DeliverCallback() {
                    @Override
                    public void handle(String consumerTag, Delivery message) throws IOException {
                        System.out.println(queueName + ",收到信息是:" + new String(message.getBody(), "UTF-8"));
                    }
                }, new CancelCallback() {
                    @Override
                    public void handle(String consumerTag) throws IOException {
    
                    }
                });
                System.out.println(queueName + ",开始接受消息");
                System.in.read();
            } catch (IOException e) {
                e.printStackTrace();
            } catch (TimeoutException e) {
                e.printStackTrace();
            } finally {
                //关闭通道和连接
                if (channel != null && channel.isOpen()) {
                    try {
                        channel.close();
                    } catch (IOException e) {
                        e.printStackTrace();
                    } catch (TimeoutException e) {
                        e.printStackTrace();
                    }
                }
    
                if (connection != null && connection.isOpen()) {
                    try {
                        connection.close();
                    } catch (IOException e) {
                        e.printStackTrace();
                    }
                }
            }
        };
    
        public static void main(String[] args) {
            new Thread(runnable, "queue1").start();
            new Thread(runnable, "queue2").start();
            new Thread(runnable, "queue3").start();
            new Thread(runnable, "queue4").start();
        }
    
    }
    
  3. 运行生产者,发送消息,运行消费者,消费消息

    queue4,收到信息是:hello direct.exchange
    queue1,收到信息是:hello direct.exchange
    

19、rabbitMQ入门案例-topic交换机

image-20230521233327970

符号 作用
. 用来分割单词
* 匹配一个单词
# 匹配一个或者多个单词
  1. 新建一个topic交换机

    image-20230521235329662

  2. 绑定队列,并制定routing key

    image-20230521235740360

  3. 发送消息

    可以发送

    image-20230522000254621

    无法匹配

    image-20230522000343509

    可以发送

    image-20230522000405453

    可以正常发送

    image-20230522000432274

    可以正常发送

    image-20230522000509300

java代码实现

  1. 生产者

    import com.rabbitmq.client.Channel;
    import com.rabbitmq.client.Connection;
    import com.rabbitmq.client.ConnectionFactory;
    
    import java.io.IOException;
    import java.util.concurrent.TimeoutException;
    
    public class Producer {
    
        public static void main(String[] args) {
            //创建连接并设置连接属性
            ConnectionFactory connectionFactory = new ConnectionFactory();
            connectionFactory.setVirtualHost("/");
            connectionFactory.setHost("39.108.49.252");
            connectionFactory.setPort(5672);
            connectionFactory.setUsername("admin");
            connectionFactory.setPassword("admin");
    
            Connection connection = null;
            Channel channel = null;
            try {
                //从连接工厂中新建连接
                connection = connectionFactory.newConnection("producer-1");
                //从连接中创建通道
                channel = connection.createChannel();
                //发送消息
                String exchangeName = "topic.exchange";
                String message = "hello topic.exchange";
                //指定了routingKey
    //            String routingKey = "1.key1.1";
    //            String routingKey = "key2.1";
    //            String routingKey = "1.key2.1";
                String routingKey = "1.1.key2.1.1";
                channel.basicPublish(exchangeName, routingKey, null, message.getBytes());
                System.out.println("消息发送成功");
            } catch (IOException e) {
                e.printStackTrace();
            } catch (TimeoutException e) {
                e.printStackTrace();
            } finally {
                //关闭通道和连接
                if (channel != null && channel.isOpen()) {
                    try {
                        channel.close();
                    } catch (IOException e) {
                        e.printStackTrace();
                    } catch (TimeoutException e) {
                        e.printStackTrace();
                    }
                }
    
                if (connection != null && connection.isOpen()) {
                    try {
                        connection.close();
                    } catch (IOException e) {
                        e.printStackTrace();
                    }
                }
            }
        }
    }
    
    
  2. 消费者

    import com.rabbitmq.client.CancelCallback;
    import com.rabbitmq.client.Channel;
    import com.rabbitmq.client.Connection;
    import com.rabbitmq.client.ConnectionFactory;
    import com.rabbitmq.client.DeliverCallback;
    import com.rabbitmq.client.Delivery;
    
    import java.io.IOException;
    import java.util.concurrent.TimeoutException;
    
    public class Consumer {
    
        private static Runnable runnable = () -> {
            //创建连接并设置连接属性
            ConnectionFactory connectionFactory = new ConnectionFactory();
            connectionFactory.setVirtualHost("/");
            connectionFactory.setHost("39.108.49.252");
            connectionFactory.setPort(5672);
            connectionFactory.setUsername("admin");
            connectionFactory.setPassword("admin");
    
            String queueName = Thread.currentThread().getName();
    
            Connection connection = null;
            Channel channel = null;
            try {
                //从连接工厂中新建连接
                connection = connectionFactory.newConnection("producer-1");
                //从连接中创建通道
                channel = connection.createChannel();
    
                //接受消息
                channel.basicConsume(queueName, true, new DeliverCallback() {
                    @Override
                    public void handle(String consumerTag, Delivery message) throws IOException {
                        System.out.println(queueName + ",收到信息是:" + new String(message.getBody(), "UTF-8"));
                    }
                }, new CancelCallback() {
                    @Override
                    public void handle(String consumerTag) throws IOException {
    
                    }
                });
                System.out.println(queueName + ",开始接受消息");
                System.in.read();
            } catch (IOException e) {
                e.printStackTrace();
            } catch (TimeoutException e) {
                e.printStackTrace();
            } finally {
                //关闭通道和连接
                if (channel != null && channel.isOpen()) {
                    try {
                        channel.close();
                    } catch (IOException e) {
                        e.printStackTrace();
                    } catch (TimeoutException e) {
                        e.printStackTrace();
                    }
                }
    
                if (connection != null && connection.isOpen()) {
                    try {
                        connection.close();
                    } catch (IOException e) {
                        e.printStackTrace();
                    }
                }
            }
        };
    
        public static void main(String[] args) {
            new Thread(runnable, "queue1").start();
            new Thread(runnable, "queue2").start();
            new Thread(runnable, "queue3").start();
            new Thread(runnable, "queue4").start();
        }
    
    }
    
    

20、rabbitMQ入门案例-headers交换机

  1. 新建一个headers交换机

    设置参数arg1=1

    image-20230522000603544

  2. 绑定队列

    image-20230522000731319

  3. 发送消息

    image-20230522000837897

    这里只有queue1收到,queue2没有收到

21、rabbitmq入门案例-work工作队列模式

image-20230522231405683

当有多个消费者时,我们的消息会被哪个消费者消费呢,我们又该如何均衡消费者消费信息的多少呢?

主要有两种模式:

  1. 轮询模式的分发:一个消费者一条,按均分配
  2. 公平分发:根据消费者的消费能力进行公平分发,处理快的处理的多,处理慢的处理的少;按劳分配;

work模式--轮询模式(round-robin)

  • 特点:该模式接收消息是当有多个消费者接入时,消息的分配模式是一个消费者分配一条,直至消息消费完成;
  1. 生产者

    import com.rabbitmq.client.Channel;
    import com.rabbitmq.client.Connection;
    import com.rabbitmq.client.ConnectionFactory;
    
    import java.io.IOException;
    import java.util.concurrent.TimeoutException;
    
    public class Producer {
    
        public static void main(String[] args) {
            ConnectionFactory connectionFactory = new ConnectionFactory();
            connectionFactory.setVirtualHost("/");
            connectionFactory.setHost("39.108.49.252");
            connectionFactory.setPort(5672);
            connectionFactory.setUsername("admin");
            connectionFactory.setPassword("admin");
    
            Connection connection = null;
            Channel channel = null;
            try {
                connection = connectionFactory.newConnection("zyy-connection");
                channel = connection.createChannel();
    
    
                for (int i = 1; i <= 20; i++) {
    
    
                    String message = "hello" + i;
    
                    channel.basicPublish("", "queue1", null, message.getBytes());
    
                }
    
                System.out.println("消息发送成功");
    
            } catch (IOException e) {
                e.printStackTrace();
            } catch (TimeoutException e) {
                e.printStackTrace();
            } finally {
                if (channel != null && channel.isOpen()) {
                    try {
                        channel.close();
                    } catch (IOException e) {
                        e.printStackTrace();
                    } catch (TimeoutException e) {
                        e.printStackTrace();
                    }
                }
                if (connection != null && connection.isOpen()) {
                    try {
                        connection.close();
                    } catch (IOException e) {
                        e.printStackTrace();
                    }
                }
    
            }
    
        }
    }
    
    
  2. 消费者1(消费能力弱)

    import com.rabbitmq.client.CancelCallback;
    import com.rabbitmq.client.Channel;
    import com.rabbitmq.client.Connection;
    import com.rabbitmq.client.ConnectionFactory;
    import com.rabbitmq.client.DeliverCallback;
    import com.rabbitmq.client.Delivery;
    
    import java.io.IOException;
    import java.util.concurrent.TimeoutException;
    
    public class Work1 {
    
        public static void main(String[] args) {
            ConnectionFactory connectionFactory = new ConnectionFactory();
            connectionFactory.setVirtualHost("/");
            connectionFactory.setHost("39.108.49.252");
            connectionFactory.setPort(5672);
            connectionFactory.setUsername("admin");
            connectionFactory.setPassword("admin");
    
            Connection connection = null;
            Channel channel = null;
            try {
                connection = connectionFactory.newConnection("zyy-connection");
                channel = connection.createChannel();
    
                channel.basicQos(1);
    
                channel.basicConsume("queue1", true, new DeliverCallback() {
                    @Override
                    public void handle(String consumerTag, Delivery message) throws IOException {
                        System.out.println("work1接收到的信息是:" + new String(message.getBody(), "UTF-8"));
                        try {
                            Thread.sleep(1000);
                        } catch (InterruptedException e) {
                            e.printStackTrace();
                        }
    
                    }
                }, new CancelCallback() {
                    @Override
                    public void handle(String consumerTag) throws IOException {
    
                    }
                });
    
                System.out.println("work1开始接受消息");
    
                System.in.read();
            } catch (IOException e) {
                e.printStackTrace();
            } catch (TimeoutException e) {
                e.printStackTrace();
            } finally {
                if (channel != null && channel.isOpen()) {
                    try {
                        channel.close();
                    } catch (IOException e) {
                        e.printStackTrace();
                    } catch (TimeoutException e) {
                        e.printStackTrace();
                    }
                }
                if (connection != null && connection.isOpen()) {
                    try {
                        connection.close();
                    } catch (IOException e) {
                        e.printStackTrace();
                    }
                }
    
            }
    
    
        }
    }
    
    
  3. 消费者2(消费能力强)

    import com.rabbitmq.client.CancelCallback;
    import com.rabbitmq.client.Channel;
    import com.rabbitmq.client.Connection;
    import com.rabbitmq.client.ConnectionFactory;
    import com.rabbitmq.client.DeliverCallback;
    import com.rabbitmq.client.Delivery;
    
    import java.io.IOException;
    import java.util.concurrent.TimeoutException;
    
    public class Work2 {
    
        public static void main(String[] args) {
            ConnectionFactory connectionFactory = new ConnectionFactory();
            connectionFactory.setVirtualHost("/");
            connectionFactory.setHost("39.108.49.252");
            connectionFactory.setPort(5672);
            connectionFactory.setUsername("admin");
            connectionFactory.setPassword("admin");
    
            Connection connection = null;
            Channel channel = null;
            try {
                connection = connectionFactory.newConnection("zyy-connection");
                channel = connection.createChannel();
    
                channel.basicQos(1);
    
    
                channel.basicConsume("queue1", true, new DeliverCallback() {
                    @Override
                    public void handle(String consumerTag, Delivery message) throws IOException {
                        System.out.println("work2接收到的信息是:" + new String(message.getBody(), "UTF-8"));
                        try {
                            Thread.sleep(100);
                        } catch (InterruptedException e) {
                            e.printStackTrace();
                        }
    
                    }
                }, new CancelCallback() {
                    @Override
                    public void handle(String consumerTag) throws IOException {
    
                    }
                });
    
                System.out.println("work2开始接受消息");
    
                System.in.read();
            } catch (IOException e) {
                e.printStackTrace();
            } catch (TimeoutException e) {
                e.printStackTrace();
            } finally {
                if (channel != null && channel.isOpen()) {
                    try {
                        channel.close();
                    } catch (IOException e) {
                        e.printStackTrace();
                    } catch (TimeoutException e) {
                        e.printStackTrace();
                    }
                }
                if (connection != null && connection.isOpen()) {
                    try {
                        connection.close();
                    } catch (IOException e) {
                        e.printStackTrace();
                    }
                }
    
            }
        }
    }
    
  4. 结果(公平分发,不论消费的能力)

    work1输出

    work1接收到的信息是:hello1
    work1接收到的信息是:hello3
    work1接收到的信息是:hello5
    work1接收到的信息是:hello7
    work1接收到的信息是:hello9
    work1接收到的信息是:hello11
    work1接收到的信息是:hello13
    work1接收到的信息是:hello15
    work1接收到的信息是:hello17
    work1接收到的信息是:hello19
    

    work2输出

    work2接收到的信息是:hello2
    work2接收到的信息是:hello4
    work2接收到的信息是:hello6
    work2接收到的信息是:hello8
    work2接收到的信息是:hello10
    work2接收到的信息是:hello12
    work2接收到的信息是:hello14
    work2接收到的信息是:hello16
    work2接收到的信息是:hello18
    work2接收到的信息是:hello20
    

work模式--公平分发(fair dispatch)

  • 特点:由于消息接收者处理消息的能力不同,存在处理快慢的问题,我们就需要能者多劳,处理快的多处理,处理慢的少处理;
  1. 生产者

    import com.rabbitmq.client.Channel;
    import com.rabbitmq.client.Connection;
    import com.rabbitmq.client.ConnectionFactory;
    
    import java.io.IOException;
    import java.util.concurrent.TimeoutException;
    
    public class Producer {
    
        public static void main(String[] args) {
            ConnectionFactory connectionFactory = new ConnectionFactory();
            connectionFactory.setVirtualHost("/");
            connectionFactory.setHost("39.108.49.252");
            connectionFactory.setPort(5672);
            connectionFactory.setUsername("admin");
            connectionFactory.setPassword("admin");
    
            Connection connection = null;
            Channel channel = null;
            try {
                connection = connectionFactory.newConnection("zyy-connection");
                channel = connection.createChannel();
    
                for (int i = 1; i <= 20; i++) {
                    String message = "hello" + i;
                    channel.basicPublish("", "queue1", null, message.getBytes());
    
                }
    
                System.out.println("消息发送成功");
    
            } catch (IOException e) {
                e.printStackTrace();
            } catch (TimeoutException e) {
                e.printStackTrace();
            } finally {
                if (channel != null && channel.isOpen()) {
                    try {
                        channel.close();
                    } catch (IOException e) {
                        e.printStackTrace();
                    } catch (TimeoutException e) {
                        e.printStackTrace();
                    }
                }
                if (connection != null && connection.isOpen()) {
                    try {
                        connection.close();
                    } catch (IOException e) {
                        e.printStackTrace();
                    }
                }
    
            }
    
        }
    }
    
  2. 消费者1(消费能力弱)

    import com.rabbitmq.client.CancelCallback;
    import com.rabbitmq.client.Channel;
    import com.rabbitmq.client.Connection;
    import com.rabbitmq.client.ConnectionFactory;
    import com.rabbitmq.client.DeliverCallback;
    import com.rabbitmq.client.Delivery;
    
    import java.io.IOException;
    import java.util.concurrent.TimeoutException;
    
    public class Work1 {
    
        public static void main(String[] args) {
            ConnectionFactory connectionFactory = new ConnectionFactory();
            connectionFactory.setVirtualHost("/");
            connectionFactory.setHost("39.108.49.252");
            connectionFactory.setPort(5672);
            connectionFactory.setUsername("admin");
            connectionFactory.setPassword("admin");
    
            Connection connection = null;
            Channel channel = null;
            try {
                connection = connectionFactory.newConnection("zyy-connection");
                channel = connection.createChannel();
    
                channel.basicQos(1);
    
                //自动应答改成手动应答
                Channel finalChannel = channel;
                channel.basicConsume("queue1", false, new DeliverCallback() {
                    @Override
                    public void handle(String consumerTag, Delivery message) throws IOException {
                        System.out.println("work1接收到的信息是:" + new String(message.getBody(), "UTF-8"));
                        finalChannel.basicAck(message.getEnvelope().getDeliveryTag(), false);
                        try {
                            Thread.sleep(1000);
                        } catch (InterruptedException e) {
                            e.printStackTrace();
                        }
    
                    }
                }, new CancelCallback() {
                    @Override
                    public void handle(String consumerTag) throws IOException {
    
                    }
                });
    
                System.out.println("work1开始接受消息");
    
                System.in.read();
            } catch (IOException e) {
                e.printStackTrace();
            } catch (TimeoutException e) {
                e.printStackTrace();
            } finally {
                if (channel != null && channel.isOpen()) {
                    try {
                        channel.close();
                    } catch (IOException e) {
                        e.printStackTrace();
                    } catch (TimeoutException e) {
                        e.printStackTrace();
                    }
                }
                if (connection != null && connection.isOpen()) {
                    try {
                        connection.close();
                    } catch (IOException e) {
                        e.printStackTrace();
                    }
                }
            }
        }
    }
    
    
  3. 消费者2(消费能力强)

    import com.rabbitmq.client.CancelCallback;
    import com.rabbitmq.client.Channel;
    import com.rabbitmq.client.Connection;
    import com.rabbitmq.client.ConnectionFactory;
    import com.rabbitmq.client.DeliverCallback;
    import com.rabbitmq.client.Delivery;
    
    import java.io.IOException;
    import java.util.concurrent.TimeoutException;
    
    public class Work2 {
    
        public static void main(String[] args) {
            ConnectionFactory connectionFactory = new ConnectionFactory();
            connectionFactory.setVirtualHost("/");
            connectionFactory.setHost("39.108.49.252");
            connectionFactory.setPort(5672);
            connectionFactory.setUsername("admin");
            connectionFactory.setPassword("admin");
    
            Connection connection = null;
            Channel channel = null;
            try {
                connection = connectionFactory.newConnection("zyy-connection");
                channel = connection.createChannel();
    
                //作用:同一时刻服务器只会发送一条消息给消费者
                channel.basicQos(1);
    
                //自动应答改成手动应答
                Channel finalChannel = channel;
                channel.basicConsume("queue1", false, new DeliverCallback() {
                    @Override
                    public void handle(String consumerTag, Delivery message) throws IOException {
                        System.out.println("work2接收到的信息是:" + new String(message.getBody(), "UTF-8"));
                        /**
                         * @param deliveryTag
                         * @param multiple
                         */
                        finalChannel.basicAck(message.getEnvelope().getDeliveryTag(), false);
                        try {
                            Thread.sleep(100);
                        } catch (InterruptedException e) {
                            e.printStackTrace();
                        }
    
                    }
                }, new CancelCallback() {
                    @Override
                    public void handle(String consumerTag) throws IOException {
    
                    }
                });
                System.out.println("work2开始接受消息");
                System.in.read();
            } catch (IOException e) {
                e.printStackTrace();
            } catch (TimeoutException e) {
                e.printStackTrace();
            } finally {
                if (channel != null && channel.isOpen()) {
                    try {
                        channel.close();
                    } catch (IOException e) {
                        e.printStackTrace();
                    } catch (TimeoutException e) {
                        e.printStackTrace();
                    }
                }
                if (connection != null && connection.isOpen()) {
                    try {
                        connection.close();
                    } catch (IOException e) {
                        e.printStackTrace();
                    }
                }
            }
        }
    }
    
  4. 结果

    work1结果

    work1接收到的信息是:hello1
    work1接收到的信息是:hello3
    work1接收到的信息是:hello14
    

    work2结果

    work2接收到的信息是:hello2
    work2接收到的信息是:hello4
    work2接收到的信息是:hello5
    work2接收到的信息是:hello6
    work2接收到的信息是:hello7
    work2接收到的信息是:hello8
    work2接收到的信息是:hello9
    work2接收到的信息是:hello10
    work2接收到的信息是:hello11
    work2接收到的信息是:hello12
    work2接收到的信息是:hello13
    work2接收到的信息是:hello15
    work2接收到的信息是:hello16
    work2接收到的信息是:hello17
    work2接收到的信息是:hello18
    work2接收到的信息是:hello19
    work2接收到的信息是:hello20
    

    可以看出来,消费能力更强的,消费的更多,多劳多得!

结论:

从结果可以看到,消费者1在相同时间内,处理了更多的消息;以上代码我们实现了公平分发模式;

  • 消费者一次接收一条消息,代码channel.BasicQos(0, 1, false);
  • 公平分发需要消费者开启手动应答,关闭自动应答
  • 关闭自动应答代码channel.BasicConsume(“queue_test”, false, consumer);
  • 消费者开启手动应答代码:channel.BasicAck(ea.DeliveryTag, false);

22、rabbitMQ使用场景

1、解耦、削峰、异步

同步异步的问题(串行)

串行方式:将订单信息写入数据库成功后,发送注册邮件,再发送注册短信。以上三个任务全部完成后,返回给客户端

image-20230524233354550

代码

public void makeOrder(){
    // 1 :保存订单 
    orderService.saveOrder();
    // 2: 发送短信服务
    messageService.sendSMS("order");//1-2 s
    // 3: 发送email服务
    emailService.sendEmail("order");//1-2 s
    // 4: 发送APP服务
    appService.sendApp("order");    
}

并行方式 异步线程池

并行方式:将订单信息写入数据库成功后,发送注册邮件的同时,发送注册短信。以上三个任务完成后,返回给客户端。与串行的差别是,并行的方式可以提高处理的时间

image-20230524233445272

代码

public void makeOrder(){
    // 1 :保存订单 
    orderService.saveOrder();
   // 相关发送
   relationMessage();
}
public void relationMessage(){
    // 异步
     theadpool.submit(new Callable<Object>{
         public Object call(){
             // 2: 发送短信服务  
             messageService.sendSMS("order");
         }
     })
    // 异步
     theadpool.submit(new Callable<Object>{
         public Object call(){
              // 3: 发送email服务
            emailService.sendEmail("order");
         }
     })
      // 异步
     theadpool.submit(new Callable<Object>{
         public Object call(){
             // 4: 发送短信服务
             appService.sendApp("order");
         }
     })
      // 异步
         theadpool.submit(new Callable<Object>{
         public Object call(){
             // 4: 发送短信服务
             appService.sendApp("order");
         }
     })
}

存在问题:

  1. 耦合度搞
  2. 需要自己写线程池,维护成本高
  3. 出现了消息可能会丢失,需要你自己做消息补偿
  4. 如何保证消息的可靠性需要自己写
  5. 如果服务器承载不了,需要自己去实现高可用

异步消息队列的方式

image-20230524233652814

好处:

  1. 完全解耦,用MQ建立桥接
  2. 有独立的线程池和运行模型
  3. 出现了消息可能会丢失,MQ有持久化功能
  4. 如何保证消息的可靠性,死信队列和消息转移的等
  5. 如果服务器承载不了,你需要自己去写高可用,HA镜像模型高可用。

按照以上约定,用户的响应时间相当于是订单信息写入数据库的时间,也就是50毫秒。注册邮件,发送短信写入消息队列后,直接返回,因此写入消息队列的速度很快,基本可以忽略,因此用户的响应时间可能是50毫秒。因此架构改变后,系统的吞吐量提高到每秒20 QPS。比串行提高了3倍,比并行提高了两倍

代码

public void makeOrder(){
    // 1 :保存订单 
    orderService.saveOrder();   
    rabbitTemplate.convertSend("ex","2","消息内容");
}

2、高内聚,低耦合

image-20230524233815877

3、流量的削峰

4、分布式事务的可靠消费和可靠生产

5、索引、缓存、静态化处理的数据同步

6、流量监控

7、日志监控(ELK)

8、下单、订单分发、抢票

23、rabbitMQ-springboot案例-fanout交换机--发布订阅模式

先新建一个maven基础模块rabbitmq-base

新建一个常量类

public class ConstantToRabbitMQ {

    public static final String EXCHANGE_NAME = "fanout_order_exchange";

    public static final String QUEUE_NAME_SMS = "sms.queue";
    public static final String QUEUE_NAME_EMAIL = "email.queue";
    public static final String QUEUE_NAME_WEIXIN = "weixin.queue";
}

生产者

  1. 新建一个springboot模板rabbitmq-springboot-producer,并添加rabbitMQ依赖

    image-20230525234553573

  2. 导包(需要依赖基本模块rabbitmq-base)

    <?xml version="1.0" encoding="UTF-8"?>
    <project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
             xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 https://maven.apache.org/xsd/maven-4.0.0.xsd">
        <modelVersion>4.0.0</modelVersion>
        <parent>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-starter-parent</artifactId>
            <version>2.4.3</version>
            <relativePath/> <!-- lookup parent from repository -->
        </parent>
        <groupId>com.zyy</groupId>
        <artifactId>rabbitmq-springboot-producer</artifactId>
        <version>0.0.1-SNAPSHOT</version>
        <name>rabbitmq-springboot-producer</name>
        <description>Demo project for Spring Boot</description>
        <properties>
            <java.version>1.8</java.version>
        </properties>
        <dependencies>
            <dependency>
                <groupId>org.springframework.boot</groupId>
                <artifactId>spring-boot-starter-amqp</artifactId>
            </dependency>
            <dependency>
                <groupId>org.springframework.boot</groupId>
                <artifactId>spring-boot-starter-web</artifactId>
            </dependency>
    
            <dependency>
                <groupId>org.springframework.boot</groupId>
                <artifactId>spring-boot-starter-test</artifactId>
                <scope>test</scope>
            </dependency>
            <dependency>
                <groupId>org.springframework.amqp</groupId>
                <artifactId>spring-rabbit-test</artifactId>
                <scope>test</scope>
            </dependency>
            <dependency>
                <groupId>com.zyy</groupId>
                <artifactId>rabbitmq-base</artifactId>
                <version>1.0-SNAPSHOT</version>
            </dependency>
        </dependencies>
    
        <build>
            <plugins>
                <plugin>
                    <groupId>org.springframework.boot</groupId>
                    <artifactId>spring-boot-maven-plugin</artifactId>
                </plugin>
            </plugins>
        </build>
    
    </project>
    
  3. 配置文件application.yaml

    server:
      port: 8081
    
    spring:
      rabbitmq:
        username: admin
        password: admin
        virtual-host: /
        host: 39.108.49.252
        port: 5672
    
  4. 配置类

    队列配置类

    import com.zyy.constants.ConstantToRabbitMQ;
    import org.springframework.amqp.core.Queue;
    import org.springframework.context.annotation.Bean;
    import org.springframework.context.annotation.Configuration;
    
    @Configuration
    public class QueueConfig {
        /**
         * 定义队列:邮件队列
         */
        @Bean("emailQueue")
        public Queue emailQueue() {
            return new Queue(ConstantToRabbitMQ.QUEUE_NAME_EMAIL, true);
        }
    
        /**
         * 定义队列:短信队列
         */
        @Bean("smsQueue")
        public Queue smsQueue() {
            return new Queue(ConstantToRabbitMQ.QUEUE_NAME_SMS, true);
        }
    
        /**
         * 定义队列:微信队列
         */
        @Bean("weixinQueue")
        public Queue weixinQueue() {
            return new Queue(ConstantToRabbitMQ.QUEUE_NAME_WEIXIN, true);
        }
    
    }
    
    

    定义交换机和绑定队列

    import com.zyy.constants.ConstantToRabbitMQ;
    import org.springframework.amqp.core.Binding;
    import org.springframework.amqp.core.BindingBuilder;
    import org.springframework.amqp.core.FanoutExchange;
    import org.springframework.amqp.core.Queue;
    import org.springframework.beans.factory.annotation.Qualifier;
    import org.springframework.context.annotation.Bean;
    import org.springframework.context.annotation.Configuration;
    
    
    @Configuration
    public class FanoutRabbitConfig {
        /**
         * 定义交换机
         */
        @Bean("fanoutExchange")
        public FanoutExchange fanoutExchange() {
            return new FanoutExchange(ConstantToRabbitMQ.EXCHANGE_NAME, true, false);
        }
    
    
        /**
         * 队列绑定交换机
         */
        @Bean
        public Binding emailBinding(@Qualifier("emailQueue") Queue queue, @Qualifier("fanoutExchange") FanoutExchange exchange) {
            return BindingBuilder.bind(queue).to(exchange);
        }
    
        @Bean
        public Binding smsBinding(@Qualifier("smsQueue") Queue queue, @Qualifier("fanoutExchange") FanoutExchange exchange) {
            return BindingBuilder.bind(queue).to(exchange);
        }
    
        @Bean
        public Binding weixinBinding(@Qualifier("weixinQueue") Queue queue, @Qualifier("fanoutExchange") FanoutExchange exchange) {
            return BindingBuilder.bind(queue).to(exchange);
        }
    }
    
  5. 发送消息

    import com.zyy.constants.ConstantToRabbitMQ;
    import org.springframework.amqp.rabbit.core.RabbitTemplate;
    import org.springframework.beans.factory.annotation.Autowired;
    import org.springframework.stereotype.Service;
    
    import java.util.UUID;
    
    
    @Service
    public class OrderService {
        @Autowired
        private RabbitTemplate rabbitTemplate;
    
        public void makeOrder(Long userId, Long productId, int num) {
            String orderNumber = UUID.randomUUID().toString();
    
            System.out.println("用户:" + userId + ",订单编码是:" + orderNumber);
            rabbitTemplate.convertAndSend(ConstantToRabbitMQ.EXCHANGE_NAME, "", orderNumber);
        }
    }
    
    
  6. 测试

    import com.zyy.producer.OrderService;
    import org.junit.jupiter.api.Test;
    import org.springframework.beans.factory.annotation.Autowired;
    import org.springframework.boot.test.context.SpringBootTest;
    
    @SpringBootTest
    class RabbitmqSpringbootProducerApplicationTests {
    
        @Autowired
        private OrderService orderService;
    
        @Test
        void contextLoads() {
            orderService.makeOrder(1L,1L,10);
        }
    }
    

    image-20230528141827902

消费者

  1. 再新建一个springboot模块rabbitmq-springboot-consumer,添加rabbitMQ依赖

  2. 导包,同生产者,这里忽略

  3. 配置文件application.yaml(就改了一下服务端口)

    server:
      port: 8080
    
    spring:
      rabbitmq:
        username: admin
        password: admin
        virtual-host: /
        host: 39.108.49.252
        port: 5672
    
  4. 定义三个消费者

    邮箱

    import com.zyy.constants.ConstantToRabbitMQ;
    import org.springframework.amqp.rabbit.annotation.RabbitListener;
    import org.springframework.stereotype.Component;
    
    @Component
    public class EmailConsumer {
        @RabbitListener(queues = ConstantToRabbitMQ.QUEUE_NAME_EMAIL)
        public void messageReceive(String message) {
            System.out.println("email------------->" + message);
        }
    }
    
    

    短信

    import com.zyy.constants.ConstantToRabbitMQ;
    import org.springframework.amqp.rabbit.annotation.RabbitListener;
    import org.springframework.stereotype.Component;
    
    
    @Component
    public class SMSConsumer {
        @RabbitListener(queues = ConstantToRabbitMQ.QUEUE_NAME_SMS)
        public void messageReceive(String message) {
            System.out.println("sms------------->" + message);
        }
    }
    
    

    微信

    import com.zyy.constants.ConstantToRabbitMQ;
    import org.springframework.amqp.rabbit.annotation.RabbitListener;
    import org.springframework.stereotype.Component;
    
    
    @Component
    public class WeixinConsumer {
        @RabbitListener(queues = ConstantToRabbitMQ.QUEUE_NAME_WEIXIN)
        public void messageReceive(String message) {
            System.out.println("weixin------------->" + message);
        }
    }
    
    
  5. 启动服务

    image-20230528142314028

24、rabbitMQ-springboot案例-direct交换机-路由模式

生产者

  1. rabbitmq-base模板的常量类ConstantToRabbitMQ新增如下

    
        public static final String EXCHANGE_NAME_DIRECT = "direct_order_exchange";
    
        public static final String QUEUE_NAME_SMS = "sms.queue";
        public static final String QUEUE_NAME_EMAIL = "email.queue";
        public static final String QUEUE_NAME_WEIXIN = "weixin.queue";
    
  2. 在rabbitmq-springboot-producer这个模块的新增配置类(定义交换机和绑定队列)

    import com.zyy.constants.ConstantToRabbitMQ;
    import org.springframework.amqp.core.Binding;
    import org.springframework.amqp.core.BindingBuilder;
    import org.springframework.amqp.core.DirectExchange;
    import org.springframework.amqp.core.Queue;
    import org.springframework.beans.factory.annotation.Qualifier;
    import org.springframework.context.annotation.Bean;
    import org.springframework.context.annotation.Configuration;
    
    @Configuration
    public class DirectRabbitConfig {
        /**
         * 定义交换机
         */
        @Bean("directExchange")
        public DirectExchange directExchange() {
            return new DirectExchange(ConstantToRabbitMQ.EXCHANGE_NAME_DIRECT, true, false);
        }
    
        /**
         * 定义队列,这里不再重复定义队列,就用
         */
    
        /**
         * 队列绑定交换机
         */
        @Bean
        public Binding emailDirectBinding(@Qualifier("emailQueue") Queue queue, @Qualifier("directExchange") DirectExchange exchange) {
            return BindingBuilder.bind(queue).to(exchange).with(ConstantToRabbitMQ.ROUTING_KEY_EMAIL);
        }
    
        @Bean
        public Binding smsDirectBinding(@Qualifier("smsQueue") Queue queue, @Qualifier("directExchange") DirectExchange exchange) {
            return BindingBuilder.bind(queue).to(exchange).with(ConstantToRabbitMQ.ROUTING_KEY_SMS);
        }
    
        @Bean
        public Binding weixinDirectBinding(@Qualifier("weixinQueue") Queue queue, @Qualifier("directExchange") DirectExchange exchange) {
            return BindingBuilder.bind(queue).to(exchange).with(ConstantToRabbitMQ.ROUTING_KEY_WEIXIN);
        }
    }
    
  3. 发送消息,在原OrderService类中添加方法

        public void makeDirectOrder(Long userId, Long productId, int num) {
            String orderNumber = UUID.randomUUID().toString();
    
            System.out.println("用户:" + userId + ",订单编码是:" + orderNumber);
    
            rabbitTemplate.convertAndSend(ConstantToRabbitMQ.EXCHANGE_NAME_DIRECT, ConstantToRabbitMQ.ROUTING_KEY_EMAIL, orderNumber);
            rabbitTemplate.convertAndSend(ConstantToRabbitMQ.EXCHANGE_NAME_DIRECT, ConstantToRabbitMQ.ROUTING_KEY_SMS, orderNumber);
    //        rabbitTemplate.convertAndSend(ConstantToRabbitMQ.EXCHANGE_NAME_DIRECT, ConstantToRabbitMQ.ROUTING_KEY_WEIXIN, orderNumber);
        }
    
  4. 测试类

        @Test
        void testDirect() {
            orderService.makeDirectOrder(2L, 2L, 10);
        }
    

    image-20230528161412272

消费者

  1. 没有任何改动,收到消息如下

    image-20230528161507047

25、rabbitMQ-springboot案例-topic交换机

生产者

  1. rabbitmq-base模板的常量类ConstantToRabbitMQ新增如下

        public static final String EXCHANGE_NAME_TOPIC = "topic_order_exchange";
    
  2. 在rabbitmq-springboot-producer这个模块的新增配置类(定义交换机和绑定队列)

    import com.zyy.constants.ConstantToRabbitMQ;
    import org.springframework.amqp.core.Binding;
    import org.springframework.amqp.core.BindingBuilder;
    import org.springframework.amqp.core.Queue;
    import org.springframework.amqp.core.TopicExchange;
    import org.springframework.beans.factory.annotation.Qualifier;
    import org.springframework.context.annotation.Bean;
    import org.springframework.context.annotation.Configuration;
    
    @Configuration
    public class TopicRabbitConfig {
        /**
         * 定义交换机
         */
        @Bean("topicExchange")
        public TopicExchange directExchange() {
            return new TopicExchange(ConstantToRabbitMQ.EXCHANGE_NAME_TOPIC, true, false);
        }
    
        /**
         * 定义队列,这里不再重复定义队列,就用
         */
    
        /**
         * 队列绑定交换机
         */
        @Bean
        public Binding emailTopicBinding(@Qualifier("emailQueue") Queue queue, @Qualifier("topicExchange") TopicExchange exchange) {
            return BindingBuilder.bind(queue).to(exchange).with("*.email.*");
        }
    
        @Bean
        public Binding smsTopicBinding(@Qualifier("smsQueue") Queue queue, @Qualifier("topicExchange") TopicExchange exchange) {
            return BindingBuilder.bind(queue).to(exchange).with("#.sms.#");
        }
    
        @Bean
        public Binding weixinTopicBinding(@Qualifier("weixinQueue") Queue queue, @Qualifier("topicExchange") TopicExchange exchange) {
            return BindingBuilder.bind(queue).to(exchange).with("weixin.#");
        }
    }
    
    
  3. 发送消息,在原OrderService类中添加方法

        public void makeDirectTopic(Long userId, Long productId, int num) {
            String orderNumber = UUID.randomUUID().toString();
    
            System.out.println("用户:" + userId + ",订单编码是:" + orderNumber);
    
    //        String routingKey = "sms.email.test";
            String routingKey = "weixin.email.test";
    
            rabbitTemplate.convertAndSend(ConstantToRabbitMQ.EXCHANGE_NAME_TOPIC, routingKey, orderNumber);
        }
    
  4. 测试

        @Test
        void testTopic() {
            orderService.makeDirectTopic(3L, 3L, 10);
        }
    

消费者

  1. 消费者不变,接收消息

    image-20230528234133265

26、rabbitMQ高级-过期时间TTL

26.1、概述

过期时间TTL表示可以对消息设置预期的时间,在这个时间内都可以被消费者接收获取;过了之后消息将自动被删除。RabbitMQ可以对消息和队列设置TTL。目前有两种方法可以设置。

  • 第一种方法是通过队列属性设置,队列中所有消息都有相同的过期时间。
  • 第二种方法是对消息进行单独设置,每条消息TTL可以不同。

如果上述两种方法同时使用,则消息的过期时间以两者之间TTL较小的那个数值为准。消息在队列的生存时间一旦超过设置的TTL值,就称为dead message被投递到死信队列, 消费者将无法再收到该消息。

26.2、设置队列TTL

  1. 在springboot案例-direct模式基本继续改造,基本模块rabbitmq-base下常量类ConstantToRabbitMQ新增如下

        public static final String ROUTING_KEY_TTL = "ttl";
        public static final String QUEUE_NAME_TTL = "ttl.queue";
    
  2. 生产者模块rabbitmq-springboot-producer,DirectRabbitConfig新增一个队列的配置

        @Bean
        public Binding ttlDirectBinding(@Qualifier("ttlQueue") Queue queue, @Qualifier("directExchange") DirectExchange exchange) {
            return BindingBuilder.bind(queue).to(exchange).with(ConstantToRabbitMQ.ROUTING_KEY_TTL);
        }
    
  3. 发送消息

        public void makeDirectOrder(Long userId, Long productId, int num) {
            String orderNumber = UUID.randomUUID().toString();
    
            System.out.println("用户:" + userId + ",订单编码是:" + orderNumber);
    
            rabbitTemplate.convertAndSend(ConstantToRabbitMQ.EXCHANGE_NAME_DIRECT, ConstantToRabbitMQ.ROUTING_KEY_TTL, orderNumber);
    //        rabbitTemplate.convertAndSend(ConstantToRabbitMQ.EXCHANGE_NAME_DIRECT, ConstantToRabbitMQ.ROUTING_KEY_SMS, orderNumber);
    //        rabbitTemplate.convertAndSend(ConstantToRabbitMQ.EXCHANGE_NAME_DIRECT, ConstantToRabbitMQ.ROUTING_KEY_WEIXIN, orderNumber);
        }
    
  4. 测试

        @Test
        void testDirect() {
            orderService.makeDirectOrder(2L, 2L, 10);
        }
    
    

    image-20230529215702855

    查看管理页面,队列创建成功,发送的一条消息5秒后就没了

    image-20230529215745634

26.3、设置消息TTL

  1. 直接改造发送消息

        public void makeDirectOrder(Long userId, Long productId, int num) {
            String orderNumber = UUID.randomUUID().toString();
    
            System.out.println("用户:" + userId + ",订单编码是:" + orderNumber);
    
            MessagePostProcessor messagePostProcessor = message -> {
                message.getMessageProperties().setExpiration("5000");
                return message;
            };
    
    
    //        rabbitTemplate.convertAndSend(ConstantToRabbitMQ.EXCHANGE_NAME_DIRECT, ConstantToRabbitMQ.ROUTING_KEY_TTL, orderNumber);
            rabbitTemplate.convertAndSend(ConstantToRabbitMQ.EXCHANGE_NAME_DIRECT, ConstantToRabbitMQ.ROUTING_KEY_SMS, orderNumber, messagePostProcessor);
    //        rabbitTemplate.convertAndSend(ConstantToRabbitMQ.EXCHANGE_NAME_DIRECT, ConstantToRabbitMQ.ROUTING_KEY_WEIXIN, orderNumber);
        }
    
  2. 测试

        @Test
        void testDirect() {
            orderService.makeDirectOrder(2L, 2L, 10);
        }
    

    image-20230529220901517

    队列收到消息后,5秒钟后消息没了

    image-20230529220941878

27、rabbitMQ高级-死信队列

概述

DLX,全称为Dead-Letter-Exchange , 可以称之为死信交换机,也有人称之为死信邮箱。当消息在一个队列中变成死信(dead message)之后,它能被重新发送到另一个交换机中,这个交换机就是DLX ,绑定DLX的队列就称之为死信队列。 消息变成死信,可能是由于以下的原因:

  • 消息被拒绝
  • 消息过期
  • 队列达到最大长度

DLX也是一个正常的交换机,和一般的交换机没有区别,它能在任何的队列上被指定,实际上就是设置某一个队列的属性。当这个队列中存在死信时,Rabbitmq就会自动地将这个消息重新发布到设置的DLX上去,进而被路由到另一个队列,即死信队列。 要想使用死信队列,只需要在定义队列的时候设置队列参数 x-dead-letter-exchange 指定交换机即可。

案例

  1. 常量类中新增如下:

        public static final String EXCHANGE_NAME_DEAD_DIRECT = "dead_direct_order_exchange";
        public static final String ROUTING_KEY_DEAD = "dead";
        public static final String QUEUE_NAME_DEAD = "dead.queue";
    
  2. 新建一个死信队列,其实就是一个普通的队列

        @Bean("deadQueue")
        public Queue deadQueue() {
            return new Queue(ConstantToRabbitMQ.QUEUE_NAME_DEAD, true);
        }
    
  3. 新建一个交换机,也就普通的交换机

        @Bean("deadExchange")
        public DirectExchange deadExchange() {
            return new DirectExchange(ConstantToRabbitMQ.EXCHANGE_NAME_DEAD_DIRECT, true, false);
        }
    
  4. 改造之前的ttlQueue队列,这里是重点,绑定的我们新建的死信队列

        @Bean("ttlQueue")
        public Queue ttlQueue() {
            //消息5秒后自动删除
            Map<String, Object> argMap = new HashMap<>();
            //过期时间,单位毫秒数
            argMap.put("x-message-ttl", 5000);
            //队列的最大长度
            argMap.put("x-max-length", 5);
            //死信队列
            argMap.put("x-dead-letter-exchange", ConstantToRabbitMQ.EXCHANGE_NAME_DEAD_DIRECT);
            //死信队列 的 routing-key
            argMap.put("x-dead-letter-routing-key", ConstantToRabbitMQ.ROUTING_KEY_DEAD);
            // (String name, boolean durable, boolean exclusive, boolean autoDelete, @Nullable Map<String, Object> arguments)
            return new Queue(ConstantToRabbitMQ.QUEUE_NAME_TTL, true, false, false, argMap);
        }
    
  5. 发送消息

        public void makeDirectOrder(Long userId, Long productId, int num) {
            String orderNumber = UUID.randomUUID().toString();
    
            System.out.println("用户:" + userId + ",订单编码是:" + orderNumber);
    
            /*MessagePostProcessor messagePostProcessor = message -> {
                MessageProperties messageProperties = message.getMessageProperties();
                messageProperties.setExpiration("5000");
                return message;
            };*/
    
            rabbitTemplate.convertAndSend(ConstantToRabbitMQ.EXCHANGE_NAME_DIRECT, ConstantToRabbitMQ.ROUTING_KEY_TTL, orderNumber);
    //        rabbitTemplate.convertAndSend(ConstantToRabbitMQ.EXCHANGE_NAME_DIRECT, ConstantToRabbitMQ.ROUTING_KEY_SMS, orderNumber, messagePostProcessor);
    //        rabbitTemplate.convertAndSend(ConstantToRabbitMQ.EXCHANGE_NAME_DIRECT, ConstantToRabbitMQ.ROUTING_KEY_WEIXIN, orderNumber);
        }
    
    
  6. 测试,发送11条信息

        @Test
        void testDirect() {
            for (int i = 0; i < 11; i++) {
                orderService.makeDirectOrder(2L, 2L, 10);
            }
        }
    

    我们设置ttlQueue队列的最大长度只有5,发送了11条消息(ps:这个时候又没有消费者),所以一开始就有6条信息直接进入死信队列中

    image-20230529224456142

    然后过了5秒后,ttl.queue队列中的5条消息都过期了,就都跑到死信队列dead.queue中了

    image-20230529224610855

28、rabbitMQ高级-内存磁盘监控

28.1、rabbitMQ的内存警告

当内存使用超过配置的阈值或者磁盘空间剩余空间小于配置的阈值时,RabbitMQ会暂时阻塞客户端的连接,并且停止接收从客户端发来的消息,以此避免服务器的崩溃,客户端与服务端的心态检测机制也会失效。

image-20230530222620641

28.2、rabbitMQ的内存控制

参考帮助文档:https://www.rabbitmq.com/configure.html 当出现警告的时候,可以通过配置去修改和调整

命令的方式

# 这里0.4表示当rabbitMQ使用的内存超过40%时,就会发送内存警告
rabbitmqctl set_vm_memory_high_watermark 0.4
# absolute参数是一个整数,表示内存使用量的绝对者,这里当rabbitMQ使用的内存超过50MB时,就会触发流量控制
rabbitmqctl set_vm_memory_high_watermark absolute 50MB

内存阈值。默认情况是:0.4,代表的含义是:当RabbitMQ的内存超过40%时,就会产生警告并且阻塞所有生产者的连接。通过此命令修改阈值在Broker重启以后将会失效,通过修改配置文件方式设置的阈值则不会随着重启而消失,但修改了配置文件一样要重启broker才会生效。

image-20230530223525763

这个时候队列都会阻塞

image-20230530223556003

配置文件方式 rabbitmq.conf

配置文件:/etc/rabbitmq/rabbitmq.conf

#默认
#vm_memory_high_watermark.relative = 0.4
# 使用relative相对值进行设置fraction,建议取值在04~0.7之间,不建议超过0.7.
vm_memory_high_watermark.relative = 0.6
# 使用absolute的绝对值的方式,但是是KB,MB,GB对应的命令如下
vm_memory_high_watermark.absolute = 2GB

28.3、rabbitMQ的内存换页

在某个Broker节点及内存阻塞生产者之前,它会尝试将队列中的消息换页到磁盘以释放内存空间,持久化和非持久化的消息都会写入磁盘中,其中持久化的消息本身就在磁盘中有一个副本,所以在转移的过程中持久化的消息会先从内存中清除掉。

默认情况下,内存到达的阈值是50%时就会换页处理。
也就是说,在默认情况下该内存的阈值是0.4的情况下,当内存超过0.4*0.5=0.2时,会进行换页动作。

比如有1000MB内存,当内存的使用率达到了400MB,已经达到了极限,但是因为配置的换页内存0.5,这个时候会在达到极限400mb之前,会把内存中的200MB进行转移到磁盘中。从而达到稳健的运行。

可以通过设置 vm_memory_high_watermark_paging_ratio 来进行调整

vm_memory_high_watermark.relative = 0.4
vm_memory_high_watermark_paging_ratio = 0.7(设置小于1的值)

为什么设置小于1,以为你如果你设置为1的阈值。内存都已经达到了极限了。你在去换页意义不是很大了。

28.4、rabbitMQ的磁盘预警

当磁盘的剩余空间低于确定的阈值时,RabbitMQ同样会阻塞生产者,这样可以避免因非持久化的消息持续换页而耗尽磁盘空间导致服务器崩溃。

默认情况下:磁盘预警为50MB的时候会进行预警。表示当前磁盘空间第50MB的时候会阻塞生产者并且停止内存消息换页到磁盘的过程。
这个阈值可以减小,但是不能完全的消除因磁盘耗尽而导致崩溃的可能性。比如在两次磁盘空间的检查空隙内,第一次检查是:60MB ,第二检查可能就是1MB,就会出现警告。

命令方式

rabbitmqctl set_disk_free_limit  <disk_limit>
rabbitmqctl set_disk_free_limit memory_limit  <fraction>
# disk_limit:固定单位 KB MB GB
# fraction :是相对阈值,建议范围在:1.0~2.0之间。(相对于内存)

配置文件方式

disk_free_limit.relative = 3.0
disk_free_limit.absolute = 50mb

29、rabbitMQ-高级-集群

29.1、rabbitMQ集群

RabbitMQ这款消息队列中间件产品本身是基于Erlang编写,Erlang([ˈɜːlæŋ])语言天生具备分布式特性(通过同步Erlang集群各节点的magic cookie来实现)。因此,RabbitMQ天然支持Clustering。这使得RabbitMQ本身不需要像ActiveMQ、Kafka那样通过ZooKeeper分别来实现HA方案和保存集群的元数据。集群是保证可靠性的一种方式,同时可以通过水平扩展以达到增加消息吞吐量能力的目的。 在实际使用过程中多采取多机多实例部署方式,为了便于同学们练习搭建,有时候你不得不在一台机器上去搭建一个rabbitmq集群,本章主要针对单机多实例这种方式来进行开展。

主要参考官方文档:https://www.rabbitmq.com/clustering.html

29.2、集群搭建

配置的前提是你的rabbitmq可以运行起来,比如”ps aux|grep rabbitmq”你能看到相关进程,又比如运行“rabbitmqctl status”你可以看到类似如下信息,而不报错:

执行下面命令进行查看:

[root@iZwz9efdd2ukk4oauustczZ ~]# ps aux|grep rabbitmq
rabbitmq 10016  0.0  0.0  48892   528 ?        S    22:07   0:00 /usr/lib64/erlang/erts-11.1.8/bin/epmd -daemon
polkitd  10728  0.0  0.0   2600   652 ?        Ss   22:21   0:00 /bin/sh /opt/rabbitmq/sbin/rabbitmq-server
polkitd  10760  0.6  7.0 2256164 133544 ?      Sl   22:21   0:25 /usr/local/lib/erlang/erts-12.2/bin/beam.smp -W w -MBas ageffcbf -MHas ageffcbf -MBlmbcs 512 -MHlmbcs 512 -MMmcs 30 -P 1048576 -t 5000000 -stbt db -zdbbl 128000 -sbwt none -sbwtdcpu none -sbwtdio none -B i -- -root /usr/local/lib/erlang -progname erl -- -home /var/lib/rabbitmq -- -pa  -noshell -noinput -s rabbit boot -boot start_sasl -syslog logger [] -syslog syslog_error_logger false
root     13782  0.0  0.0 112816   980 pts/2    S+   23:30   0:00 grep --color=auto rabbitmq
[root@iZwz9efdd2ukk4oauustczZ ~]# systemctl status rabbitmq-server # 或者这个命令
● rabbitmq-server.service - RabbitMQ broker
   Loaded: loaded (/usr/lib/systemd/system/rabbitmq-server.service; disabled; vendor preset: disabled)
   Active: inactive (dead)

May 17 17:44:50 iZwz9efdd2ukk4oauustczZ rabbitmq-server[828]: /var/log/rabbitmq/rabbit@iZwz9efdd2ukk4oauustczZ_upgrade.log
May 17 17:44:50 iZwz9efdd2ukk4oauustczZ rabbitmq-server[828]: Config file(s): (none)
May 17 17:44:52 iZwz9efdd2ukk4oauustczZ rabbitmq-server[828]: Starting broker... completed with 3 plugins.
May 17 17:44:52 iZwz9efdd2ukk4oauustczZ systemd[1]: Started RabbitMQ broker.
May 17 21:59:00 iZwz9efdd2ukk4oauustczZ systemd[1]: Stopping RabbitMQ broker...
May 17 21:59:01 iZwz9efdd2ukk4oauustczZ rabbitmqctl[11679]: Shutting down RabbitMQ node rabbit@iZwz9efdd2ukk4oauustczZ running at PID 828
May 17 21:59:01 iZwz9efdd2ukk4oauustczZ rabbitmq-server[828]: Gracefully halting Erlang VM
May 17 21:59:02 iZwz9efdd2ukk4oauustczZ rabbitmqctl[11679]: Waiting for PID 828 to terminate
May 17 21:59:09 iZwz9efdd2ukk4oauustczZ rabbitmqctl[11679]: RabbitMQ node rabbit@iZwz9efdd2ukk4oauustczZ running at PID 828 successfully shut down
May 17 21:59:09 iZwz9efdd2ukk4oauustczZ systemd[1]: Stopped RabbitMQ broker.
[root@iZwz9efdd2ukk4oauustczZ ~]# 

注意:确保RabbitMQ可以运行的,确保完成之后,把单机版的RabbitMQ服务停止,后台看不到RabbitMQ的进程为止

29.3、单机多实例搭建

**场景:**假设有两个rabbitmq节点,分别为rabbit-1, rabbit-2,rabbit-1作为主节点,rabbit-2作为从节点。

启动命令:RABBITMQ_NODE_PORT=5672 RABBITMQ_NODENAME=rabbit-1 rabbitmq-server -detached

结束命令:rabbitmqctl -n rabbit-1 stop

第一步:启动第一个节点rabbit-1

[root@iZwz9efdd2ukk4oauustczZ ~]# RABBITMQ_NODE_PORT=5672 RABBITMQ_NODENAME=rabbit-1 rabbitmq-server start
Configuring logger redirection

  ##  ##      RabbitMQ 3.8.12
  ##  ##
  ##########  Copyright (c) 2007-2021 VMware, Inc. or its affiliates.
  ######  ##
  ##########  Licensed under the MPL 2.0. Website: https://rabbitmq.com

  Doc guides: https://rabbitmq.com/documentation.html
  Support:    https://rabbitmq.com/contact.html
  Tutorials:  https://rabbitmq.com/getstarted.html
  Monitoring: https://rabbitmq.com/monitoring.html

  Logs: /var/log/rabbitmq/rabbit-1@iZwz9efdd2ukk4oauustczZ.log
        /var/log/rabbitmq/rabbit-1@iZwz9efdd2ukk4oauustczZ_upgrade.log

  Config file(s): (none)

  Starting broker... completed with 3 plugins.

至此节点rabbit-1启动完成

第二步:启动第二个节点rabbit-2

注意:web管理插件端口占用,所以还要指定其web插件占用的端口号 RABBITMQ_SERVER_START_ARGS=”-rabbitmq_management listener [{port,15673}]”

[root@iZwz9efdd2ukk4oauustczZ ~]# RABBITMQ_NODE_PORT=5673 RABBITMQ_SERVER_START_ARGS="-rabbitmq_management listener [{port,15673}]" RABBITMQ_NODENAME=rabbit-2 rabbitmq-server start
Configuring logger redirection

  ##  ##      RabbitMQ 3.8.12
  ##  ##
  ##########  Copyright (c) 2007-2021 VMware, Inc. or its affiliates.
  ######  ##
  ##########  Licensed under the MPL 2.0. Website: https://rabbitmq.com

  Doc guides: https://rabbitmq.com/documentation.html
  Support:    https://rabbitmq.com/contact.html
  Tutorials:  https://rabbitmq.com/getstarted.html
  Monitoring: https://rabbitmq.com/monitoring.html

  Logs: /var/log/rabbitmq/rabbit-2@iZwz9efdd2ukk4oauustczZ.log
        /var/log/rabbitmq/rabbit-2@iZwz9efdd2ukk4oauustczZ_upgrade.log

  Config file(s): (none)

  Starting broker... completed with 3 plugins.

至此节点rabbit-2启动完成

第三步:验证启动ps aux|grep rabbitmq

image-20230530234201836

第四步:rabbit-1操作作为主节点

# 停止应用
[root@iZwz9efdd2ukk4oauustczZ ~]# rabbitmqctl -n rabbit-1 stop_app 
Stopping rabbit application on node rabbit-1@iZwz9efdd2ukk4oauustczZ ...
# 目的是清除节点上的历史数据(如果不清除,无法将节点加入到集群)
[root@iZwz9efdd2ukk4oauustczZ ~]# rabbitmqctl -n rabbit-1 reset 
Resetting node rabbit-1@iZwz9efdd2ukk4oauustczZ ...
#启动应用
[root@iZwz9efdd2ukk4oauustczZ ~]# rabbitmqctl -n rabbit-1 start_app 
Starting node rabbit-1@iZwz9efdd2ukk4oauustczZ ...
[root@iZwz9efdd2ukk4oauustczZ ~]# 

第五步:rabbit-2操作作为从节点

# 停止应用
[root@iZwz9efdd2ukk4oauustczZ ~]# rabbitmqctl -n rabbit-2 stop_app
Stopping rabbit application on node rabbit-2@iZwz9efdd2ukk4oauustczZ ...
# 目的是清除节点上的历史数据(如果不清除,无法将节点加入到集群)
[root@iZwz9efdd2ukk4oauustczZ ~]# rabbitmqctl -n rabbit-2 reset
Resetting node rabbit-2@iZwz9efdd2ukk4oauustczZ ...
# 将rabbit2节点加入到rabbit1(主节点)集群当中【Server-node服务器的主机名】
# rabbitmqctl -n rabbit-2 join_cluster rabbit-1@'Server-node'
[root@iZwz9efdd2ukk4oauustczZ ~]# rabbitmqctl -n rabbit-2 join_cluster rabbit-1@iZwz9efdd2ukk4oauustczZ
Clustering node rabbit-2@iZwz9efdd2ukk4oauustczZ with rabbit-1@iZwz9efdd2ukk4oauustczZ
# 启动应用
[root@iZwz9efdd2ukk4oauustczZ ~]# rabbitmqctl -n rabbit-2 start_app
Starting node rabbit-2@iZwz9efdd2ukk4oauustczZ ...
[root@iZwz9efdd2ukk4oauustczZ ~]# 

第六步:验证集群状态

[root@iZwz9efdd2ukk4oauustczZ ~]# rabbitmqctl cluster_status -n rabbit-1

image-20230530234822685

第七步:web监控

执行命令

[root@iZwz9efdd2ukk4oauustczZ ~]# rabbitmq-plugins enable rabbitmq_management
Enabling plugins on node rabbit@iZwz9efdd2ukk4oauustczZ:
rabbitmq_management
The following plugins have been configured:
  rabbitmq_management
  rabbitmq_management_agent
  rabbitmq_web_dispatch
Applying plugin configuration to rabbit@iZwz9efdd2ukk4oauustczZ...
Plugin configuration unchanged.
[root@iZwz9efdd2ukk4oauustczZ ~]# 

注意在访问的时候:web结面的管理需要给15672 node-1 和15673的node-2 设置用户名和密码。如下:

rabbitmqctl -n rabbit-1 add_user admin admin
rabbitmqctl -n rabbit-1 set_user_tags admin administrator
rabbitmqctl -n rabbit-1 set_permissions -p / admin ".*" ".*" ".*"
rabbitmqctl -n rabbit-2 add_user admin admin
rabbitmqctl -n rabbit-2 set_user_tags admin administrator
rabbitmqctl -n rabbit-2 set_permissions -p / admin ".*" ".*" ".*"

image-20230607235604434

小结

Tips:
如果采用多机部署方式,需读取其中一个节点的cookie, 并复制到其他节点(节点之间通过cookie确定相互是否可通信)。cookie存放在/var/lib/rabbitmq/.erlang.cookie。
例如:主机名分别为rabbit-1、rabbit-2
1、逐个启动各节点
2、配置各节点的hosts文件( vim /etc/hosts)
​ ip1:rabbit-1
​ ip2:rabbit-2
其它步骤雷同单机部署方式

30、rabbitMQ-高级-分布式事务

简介

分布式事务指事务的操作位于不同的节点上,需要保证事务的 AICD 特性。 例如在下单场景下,库存和订单如果不在同一个节点上,就涉及分布式事务。

30.1、分布式事务的方式

在分布式系统中,要实现分布式事务,无外乎那几种解决方案。

1、两阶段提交(2PC)需要数据库产商的支持,java组件有atomikos等

两阶段提交(Two-phase Commit,2PC),通过引入协调者(Coordinator)来协调参与者的行为,并最终决定这些参与者是否要真正执行事务

准备阶段

协调者询问参与者事务是否执行成功,参与者发回事务执行结果。

image-20230601235548328

提交阶段

如果事务在每个参与者上都执行成功,事务协调者发送通知让参与者提交事务;否则,协调者发送通知让参与者回滚事务。 需要注意的是,在准备阶段,参与者执行了事务,但是还未提交。只有在提交阶段接收到协调者发来的通知后,才进行提交或者回滚

image-20230601235652883

存在的问题

  • 同步阻塞 所有事务参与者在等待其它参与者响应的时候都处于同步阻塞状态,无法进行其它操作。
  • 单点问题 协调者在 2PC 中起到非常大的作用,发生故障将会造成很大影响。特别是在阶段二发生故障,所有参与者会一直等待状态,无法完成其它操作
  • 数据不一致 在阶段二,如果协调者只发送了部分 Commit 消息,此时网络发生异常,那么只有部分参与者接收到 Commit 消息,也就是说只有部分参与者提交了事务,使得系统数据不一致
  • 太过保守 任意一个节点失败就会导致整个事务失败,没有完善的容错机制

2、补偿事务(TCC)严选,阿里,蚂蚁金服

TCC 其实就是采用的补偿机制,其核心思想是:针对每个操作,都要注册一个与其对应的确认和补偿(撤销)操作。它分为三个阶段:

  • Try 阶段主要是对业务系统做检测及资源预留
  • Confirm 阶段主要是对业务系统做确认提交,Try阶段执行成功并开始执行 Confirm阶段时,默认 - - - Confirm阶段是不会出错的。即:只要Try成功,Confirm一定成功
  • Cancel 阶段主要是在业务执行错误,需要回滚的状态下执行的业务取消,预留资源释放
举个例子,假入 Bob 要向 Smith 转账,思路大概是: 我们有一个本地方法,里面依次调用
1:首先在 Try 阶段,要先调用远程接口把 Smith 和 Bob 的钱给冻结起来。
2:在 Confirm 阶段,执行远程调用的转账的操作,转账成功进行解冻。
3:如果第2步执行成功,那么转账成功,如果第二步执行失败,则调用远程冻结接口对应的解冻方法 (Cancel)

优点

  • 跟2PC比起来,实现以及流程相对简单了一些,但数据的一致性比2PC也要差一些

缺点:

  • 缺点还是比较明显的,在2,3步中都有可能失败。TCC属于应用层的一种补偿方式,所以需要程序员在实现的时候多写很多补偿的代码,在一些场景中,一些业务流程可能用TCC不太好定义及处理。

3、本次消息表(异步确保)比如:支付宝、微信支付主动查询支付状态,对账单的形式

本地消息表与业务数据表处于同一个数据库中,这样就能利用本地事务来保证在对这两个表的操作满足事务特性,并且使用了消息队列来保证最终一致性。

  • 在分布式事务操作的一方完成写业务数据的操作之后向本地消息表发送一个消息,本地事务能保证这个消息一定会被写入本地消息表中。
  • 之后将本地消息表中的消息转发到 Kafka 等消息队列中,如果转发成功则将消息从本地消息表中删除,否则继续重新转发
  • 在分布式事务操作的另一方从消息队列中读取一个消息,并执行消息中的操作

image-20230602000323799

优点

  • 一种非常经典的实现,避免了分布式事务,实现了最终一致性

缺点

  • 消息表会耦合到业务系统中,如果没有封装好的解决方案,会有很多杂活需要处理。

4、MQ事务消息异步场景,通用性较强,拓展性较高

有一些第三方的MQ是支持事务消息的,比如RocketMQ,他们支持事务消息的方式也是类似于采用的二阶段提交,但是市面上一些主流的MQ都是不支持事务消息的,比如 Kafka 不支持。 以阿里的 RocketMQ中间件为例,其思路大致为:

  • 第一阶段Prepared消息,会拿到消息的地址。 第二阶段执行本地事务,第三阶段通过第一阶段拿到的地址去访问消息,并修改状态
  • 也就是说在业务方法内要想消息队列提交两次请求,一次发送消息和一次确认消息。如果确认消息发送失败了RocketMQ会定期扫描消息集群中的事务消息,这时候发现了Prepared消息,它会向消息发送者确认,所以生产方需要实现一个check接口,RocketMQ会根据发送端设置的策略来决定是回滚还是继续发送确认消息。这样就保证了消息发送与本地事务同时成功或同时失败

image-20230602000526153

优点

  • 实现了最终一致性,不需要依赖本地数据库事务

缺点

  • 实现难度大,主流MQ不支持,RocketMQ事务消息部分代码也未开源

5、总结

通过本文我们总结并对比了几种分布式分解方案的优缺点,分布式事务本身是一个技术难题,是没有一种完美的方案应对所有场景的,具体还是要根据业务场景去抉择吧。阿里RocketMQ去实现的分布式事务,现在也有除了很多分布式事务的协调器,比如LCN等,大家可以多去尝试。

30.2、具体实现

数据库+表

-- 创建数据库
CREATE DATABASE `rabbitmq-dispatcher`  DEFAULT CHARACTER SET utf8 ;

-- 使用数据库
USE `rabbitmq-dispatcher`;

-- 创建表
DROP TABLE IF EXISTS `dispatcher`;

CREATE TABLE `dispatcher` (
  `dispatcher_id` BIGINT(20) NOT NULL AUTO_INCREMENT COMMENT '配送表id',
  `order_id` BIGINT(20) DEFAULT NULL COMMENT '订单id',
  `status` TINYINT DEFAULT NULL COMMENT '配送状态',
  `create_time` DATETIME DEFAULT NULL COMMENT '创建时间',
  PRIMARY KEY (`dispatcher_id`)
) ENGINE=INNODB AUTO_INCREMENT=1 DEFAULT CHARSET=utf8 COMMENT='配送表';
-- 创建数据库
CREATE DATABASE `rabbitmq-order`  DEFAULT CHARACTER SET utf8 ;

-- 使用数据库
USE `rabbitmq-order`;

-- 创建表
DROP TABLE IF EXISTS `order`;

CREATE TABLE `order` (
  `order_id` BIGINT(20) NOT NULL AUTO_INCREMENT COMMENT '订单id',
  `user_id` BIGINT(20) DEFAULT NULL COMMENT '用户id',
  `order_content` TEXT DEFAULT NULL COMMENT '订单内容',
  `create_time` DATETIME DEFAULT NULL COMMENT '创建时间',
  PRIMARY KEY (`order_id`)
) ENGINE=INNODB AUTO_INCREMENT=1 DEFAULT CHARSET=utf8 COMMENT='订单表';

image-20230603173540197

代码

代码:https://gitee.com/zhayuyao/my-rabbitmq

image-20230605203445765

配送服务

  1. 新建订单服务模块rabbit-dispatcher-service

  2. 导包

    <?xml version="1.0" encoding="UTF-8"?>
    <project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
             xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 https://maven.apache.org/xsd/maven-4.0.0.xsd">
        <modelVersion>4.0.0</modelVersion>
        <parent>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-starter-parent</artifactId>
            <version>2.4.3</version>
            <relativePath/> <!-- lookup parent from repository -->
        </parent>
        <groupId>com.zyy</groupId>
        <artifactId>rabbit-dispatcher-service</artifactId>
        <version>0.0.1-SNAPSHOT</version>
        <name>rabbit-dispatcher-service</name>
        <description>Demo project for Spring Boot</description>
        <properties>
            <java.version>1.8</java.version>
        </properties>
        <dependencies>
            <dependency>
                <groupId>org.springframework.boot</groupId>
                <artifactId>spring-boot-starter-amqp</artifactId>
            </dependency>
            <dependency>
                <groupId>org.springframework.boot</groupId>
                <artifactId>spring-boot-starter-web</artifactId>
            </dependency>
    
            <dependency>
                <groupId>org.springframework.boot</groupId>
                <artifactId>spring-boot-starter-jdbc</artifactId>
            </dependency>
    
            <dependency>
                <groupId>mysql</groupId>
                <artifactId>mysql-connector-java</artifactId>
                <scope>runtime</scope>
            </dependency>
    
            <dependency>
                <groupId>org.projectlombok</groupId>
                <artifactId>lombok</artifactId>
            </dependency>
    
            <dependency>
                <groupId>org.springframework.boot</groupId>
                <artifactId>spring-boot-starter-test</artifactId>
                <scope>test</scope>
            </dependency>
            <dependency>
                <groupId>org.springframework.amqp</groupId>
                <artifactId>spring-rabbit-test</artifactId>
                <scope>test</scope>
            </dependency>
        </dependencies>
    
        <build>
            <plugins>
                <plugin>
                    <groupId>org.springframework.boot</groupId>
                    <artifactId>spring-boot-maven-plugin</artifactId>
                </plugin>
            </plugins>
        </build>
    
    </project>
    
  3. 配置文件

    server:
      port: 9000
    
    spring:
      datasource:
        url: jdbc:mysql://localhost:3306/rabbitmq-dispatcher?useUnicode=true&characterEncoding=utf-8
        username: root
        password: 123456
        driver-class-name: com.mysql.cj.jdbc.Driver
    
      rabbitmq:
        username: admin
        password: admin
        virtual-host: /
        host: 39.108.49.252
        port: 5672
    
  4. 新建service

    import org.springframework.beans.factory.annotation.Autowired;
    import org.springframework.jdbc.core.JdbcTemplate;
    import org.springframework.stereotype.Service;
    import org.springframework.transaction.annotation.Transactional;
    
    import java.util.Date;
    
    @Service
    public class DispatcherService {
    
        @Autowired
        private JdbcTemplate jdbcTemplate;
    
        @Transactional(rollbackFor = Exception.class)
        public void createDispatcher(int orderId) throws Exception {
    
            String sql = "insert into dispatcher(order_id, status, create_time) values(?,?,?)";
    
            Date date = new Date();
            int count = jdbcTemplate.update(sql, orderId, 0, date);
            if (count != 1) {
                throw new Exception("配送单创建失败!");
            }
        }
    }
    
    
  5. 新建controller

    import com.zyy.rabbitmq.service.DispatcherService;
    import org.springframework.beans.factory.annotation.Autowired;
    import org.springframework.web.bind.annotation.GetMapping;
    import org.springframework.web.bind.annotation.RequestParam;
    import org.springframework.web.bind.annotation.RestController;
    
    @RestController
    public class DispatcherController {
        @Autowired
        DispatcherService dispatcherService;
    
        /**
         * 创建配送单
         *
         * @param orderId
         */
        @GetMapping("/dispatcher/create")
        public String createDispatcher(@RequestParam Integer orderId) {
            try {
                if (orderId.equals(8)) {
                    //模拟耗时,接口调用者会认为超时
                    Thread.sleep(4000L);
                }
                dispatcherService.createDispatcher(orderId);
                return "success";
            } catch (Exception e) {
                e.printStackTrace();
            }
            return "";
    
        }
    }
    
    

订单服务

  1. 新建配送服务模块rabbitmq-order-service

  2. 导包

    <?xml version="1.0" encoding="UTF-8"?>
    <project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
             xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 https://maven.apache.org/xsd/maven-4.0.0.xsd">
        <modelVersion>4.0.0</modelVersion>
        <parent>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-starter-parent</artifactId>
            <version>2.4.3</version>
            <relativePath/> <!-- lookup parent from repository -->
        </parent>
        <groupId>com.zyy</groupId>
        <artifactId>rabbitmq-order-service</artifactId>
        <version>0.0.1-SNAPSHOT</version>
        <name>rabbitmq-order-service</name>
        <description>Demo project for Spring Boot</description>
        <properties>
            <java.version>1.8</java.version>
        </properties>
        <dependencies>
            <dependency>
                <groupId>org.springframework.boot</groupId>
                <artifactId>spring-boot-starter-amqp</artifactId>
            </dependency>
            <dependency>
                <groupId>org.springframework.boot</groupId>
                <artifactId>spring-boot-starter-web</artifactId>
            </dependency>
            <dependency>
                <groupId>org.springframework.boot</groupId>
                <artifactId>spring-boot-starter-jdbc</artifactId>
            </dependency>
    
            <dependency>
                <groupId>mysql</groupId>
                <artifactId>mysql-connector-java</artifactId>
                <scope>runtime</scope>
            </dependency>
            <dependency>
                <groupId>org.projectlombok</groupId>
                <artifactId>lombok</artifactId>
            </dependency>
            <dependency>
                <groupId>org.springframework.boot</groupId>
                <artifactId>spring-boot-starter-test</artifactId>
                <scope>test</scope>
            </dependency>
            <dependency>
                <groupId>org.springframework.amqp</groupId>
                <artifactId>spring-rabbit-test</artifactId>
                <scope>test</scope>
            </dependency>
        </dependencies>
    
        <build>
            <plugins>
                <plugin>
                    <groupId>org.springframework.boot</groupId>
                    <artifactId>spring-boot-maven-plugin</artifactId>
                </plugin>
            </plugins>
        </build>
    
    </project>
    
    
  3. 配置文件

    server:
      port: 9001
    
    spring:
      datasource:
        url: jdbc:mysql://localhost:3306/rabbitmq-order?useUnicode=true&characterEncoding=utf-8
        username: root
        password: 123456
        driver-class-name: com.mysql.cj.jdbc.Driver
    
      rabbitmq:
        username: admin
        password: admin
        virtual-host: /
        host: 39.108.49.252
        port: 5672
    
  4. 新建实体类

    import lombok.Data;
    
    import java.io.Serializable;
    import java.util.Date;
    
    @Data
    public class Order implements Serializable {
        private static final long serialVersionUID = 2231174433174810752L;
    
        public Integer orderId;
        public Integer userId;
        public String orderContent;
        public Date createTime;
    
    }
    
    
  5. 新建service

    import com.zyy.rabbitmq.pojo.Order;
    import org.springframework.beans.factory.annotation.Autowired;
    import org.springframework.boot.web.client.RestTemplateBuilder;
    import org.springframework.http.client.SimpleClientHttpRequestFactory;
    import org.springframework.jdbc.core.JdbcTemplate;
    import org.springframework.stereotype.Service;
    import org.springframework.transaction.annotation.Transactional;
    import org.springframework.web.client.RestTemplate;
    
    
    @Service
    public class OrderService {
    
        @Autowired
        private JdbcTemplate jdbcTemplate;
    
        @Transactional(rollbackFor = Exception.class)
        public void createOrder(Order order) throws Exception {
            //先创建订单
            saveOrder(order);
    
            //然后调用远程服务,创建运单
            String result = dispatcherHttpApi(order.getOrderId());
            if (!"success".equals(result)) {
                throw new Exception("创建运单失败");
            }
    
    
        }
    
        public void saveOrder(Order order) throws Exception {
            String sql = "insert into `order`(order_id, user_id, order_content, create_time) values(?,?,?,?)";
    
            int count = jdbcTemplate.update(sql, order.getOrderId(),
                    order.getUserId(), order.getOrderContent(), order.getCreateTime());
            if (count != 1) {
                throw new Exception("订单创建失败!");
            }
    
        }
    
        private String dispatcherHttpApi(Integer orderId) {
            SimpleClientHttpRequestFactory httpRequestFactory = new SimpleClientHttpRequestFactory();
            //连接超时时间  milliseconds
            httpRequestFactory.setConnectTimeout(3000);
            //处理超时时间
            httpRequestFactory.setReadTimeout(2000);
    
            RestTemplate restTemplate = new RestTemplate(httpRequestFactory);
    
            String url = "http://localhost:9000/dispatcher/create?orderId=" + orderId;
    
            return restTemplate.getForObject(url, String.class);
    
        }
    }
    
    
  6. 测试类

    import com.zyy.rabbitmq.pojo.Order;
    import com.zyy.rabbitmq.service.OrderService;
    import org.junit.jupiter.api.Test;
    import org.springframework.beans.factory.annotation.Autowired;
    import org.springframework.boot.test.context.SpringBootTest;
    
    import java.util.Date;
    
    
    @SpringBootTest
    class RabbitmqOrderServiceApplicationTests {
    
    
        @Autowired
        private OrderService orderService;
    
        @Test
        void orderCreate() {
            Order order = new Order();
            order.setUserId(1);
            order.setOrderId(8);
            order.setOrderContent("重庆鸡公煲");
            order.setCreateTime(new Date());
    
            try {
                orderService.createOrder(order);
            } catch (Exception e) {
                e.printStackTrace();
            }
    
        }
    }
    

    启动配送服务,运行订单服务中的测试类

    正常情况,订单表会有一个数据,配送表会有一条对应的数据

    当配送表运行超过订单服务指定的时间,订单事务回滚,订单表无数据,但是配送表是插入成功的,这种就导致不一致了!

改成MQ的分布式事务处理

image-20230605203627678

如何保证可靠生产

image-20230605203953572

代码实现

  1. 在新建一个订单冗余表

    -- 使用数据库
    USE `rabbitmq-order`;
    
    -- 创建表
    DROP TABLE IF EXISTS `order_message`;
    
    CREATE TABLE `order_message` (
      `order_id` BIGINT(20) NOT NULL AUTO_INCREMENT COMMENT '订单id',
      `user_id` BIGINT(20) DEFAULT NULL COMMENT '用户id',
      `order_content` TEXT DEFAULT NULL COMMENT '订单内容',
      `status` TINYINT DEFAULT NULL COMMENT '0-处理中 1-推送成功 2-推送失败',
      `create_time` DATETIME DEFAULT NULL COMMENT '创建时间',
      PRIMARY KEY (`order_id`)
    ) ENGINE=INNODB AUTO_INCREMENT=1 DEFAULT CHARSET=utf8 COMMENT='订单冗余表';
    
  2. 导入json包

    
            <dependency>
                <groupId>com.alibaba.fastjson2</groupId>
                <artifactId>fastjson2</artifactId>
                <version>2.0.28</version>
            </dependency>
    
  3. 入订单表,入冗余表,发送消息OrderService添加方法

        @Autowired
        private JdbcTemplate jdbcTemplate;
        @Autowired
        private MQOrderService mqOrderService;
    
        @Transactional(rollbackFor = Exception.class)
        public void createOrderMessage(Order order) throws Exception {
            //先创建订单 并且冗余一份
            saveOrder(order);
            saveOrderMessage(order);
            //mq推送给运单
            mqOrderService.sendMessage(order);
        }
    
        public void saveOrder(Order order) throws Exception {
            String sql = "insert into `order`(order_id, user_id, order_content, create_time) values(?,?,?,?)";
    
            int count = jdbcTemplate.update(sql, order.getOrderId(),
                    order.getUserId(), order.getOrderContent(), order.getCreateTime());
            if (count != 1) {
                throw new Exception("订单创建失败!");
            }
        }
    
        public void saveOrderMessage(Order order) throws Exception {
            String sql = "insert into `order_message`(order_id, user_id, order_content, status, create_time) values(?,?,?,?,?)";
    
            int count = jdbcTemplate.update(sql, order.getOrderId(),
                    order.getUserId(), order.getOrderContent(), "0", order.getCreateTime());
            if (count != 1) {
                throw new Exception("订单冗余失败!");
            }
        }
    
    

    新增消息发送类

    import com.alibaba.fastjson2.JSON;
    import com.zyy.rabbitmq.pojo.Order;
    import org.springframework.amqp.rabbit.connection.CorrelationData;
    import org.springframework.amqp.rabbit.core.RabbitTemplate;
    import org.springframework.beans.factory.annotation.Autowired;
    import org.springframework.jdbc.core.JdbcTemplate;
    import org.springframework.stereotype.Service;
    
    import javax.annotation.PostConstruct;
    
    @Service
    public class MQOrderService {
    
        @Autowired
        private RabbitTemplate rabbitTemplate;
    
        @Autowired
        private JdbcTemplate jdbcTemplate;
    
        @PostConstruct
        public void regCallBack() {
            rabbitTemplate.setConfirmCallback((correlationData, ack, s) -> {
                String id = correlationData.getId();
                Integer orderId = Integer.parseInt(id);
                if (ack) {
                    updateOrderMessage(1,orderId);
                } else {
                    updateOrderMessage(2,orderId);
                }
            });
        }
    
        public void updateOrderMessage(Integer status, Integer orderId)  {
            //TODO 注意 失败情况,重试次数应该+1
            String sql = "update order_message set status=? where order_id=?";
    
            int count = jdbcTemplate.update(sql, status, orderId);
            if (count != 1) {
                //监控
            }
        }
    
        public void sendMessage(Order order) {
            String message = JSON.toJSONString(order);
            rabbitTemplate.convertAndSend("order.fanout.exchange", "", message, new CorrelationData(order.getOrderId() + ""));
        }
    }
    
    
  4. 页面上添加交换机,新增队列,然后绑定到交换机上

    image-20230608234640210

    image-20230608234656995

    image-20230608234718870

  5. 测试

        @Test
        void orderCreateMessage() {
            Order order = new Order();
            order.setUserId(1);
            order.setOrderId(1000);
            order.setOrderContent("重庆鸡公煲");
            order.setCreateTime(new Date());
    
            try {
                orderService.createOrderMessage(order);
            } catch (Exception e) {
                e.printStackTrace();
            }
    
        }
    

    然后,订单表和冗余表都会有记录,也有有一条消息

    image-20230608234919637

    image-20230608234931934

    image-20230608234958355

    后面如果生产者发送消息失败,那么冗余表就会有条状态为失败的数据,写个定时任务扫描失败数据,重新推送,需要注意控制重试次数,这样就可以保证生产者相对可靠生产

    启动类上开启定时任务

    import org.springframework.boot.SpringApplication;
    import org.springframework.boot.autoconfigure.SpringBootApplication;
    import org.springframework.scheduling.annotation.EnableScheduling;
    
    @SpringBootApplication
    @EnableScheduling
    public class RabbitmqOrderServiceApplication {
    
        public static void main(String[] args) {
            SpringApplication.run(RabbitmqOrderServiceApplication.class, args);
        }
    
    }
    

    定时任务

    import com.alibaba.fastjson2.JSON;
    import com.zyy.rabbitmq.pojo.Order;
    import com.zyy.rabbitmq.service.MQOrderService;
    import org.springframework.jdbc.core.JdbcTemplate;
    import org.springframework.scheduling.annotation.Scheduled;
    import org.springframework.stereotype.Component;
    
    import javax.annotation.Resource;
    import java.util.List;
    import java.util.Map;
    
    @Component
    public class RetryTask {
        @Resource
        private JdbcTemplate jdbcTemplate;
        @Resource
        private MQOrderService mqOrderService;
    
        @Scheduled(cron = "*/10 * * * * ?")
        public void retry() {
            //查询发送失败的任务  TODO 注意:这里应该需要控制重试次数,不能无限制重试
            String sql = "SELECT order_id 'orderId',user_id 'userId',order_content 'orderContent',`status`,create_time 'createTime' FROM order_message WHERE STATUS=2";
    
            List<Map<String, Object>> taskList = jdbcTemplate.queryForList(sql);
            if (taskList == null || taskList.size() == 0) {
                return;
            }
            List<Order> orderList = JSON.parseArray(JSON.toJSONString(taskList), Order.class);
            for (Order order : orderList) {
                //mq推送给运单
                mqOrderService.sendMessage(order);
            }
        }
    }
    
    

如何保证可靠消费

image-20230605204118858

代码实现

  1. 把消费者先改成mq

    先添加配置

      rabbitmq:
        username: admin
        password: admin
        virtual-host: /
        #    host: 39.108.49.252
        #    port: 5672
        addresses: 39.108.49.252:5672
    

    新增消费类

    import com.alibaba.fastjson2.JSON;
    import com.alibaba.fastjson2.JSONObject;
    import com.zyy.rabbitmq.service.DispatcherService;
    import org.springframework.amqp.rabbit.annotation.RabbitListener;
    import org.springframework.beans.factory.annotation.Autowired;
    import org.springframework.stereotype.Component;
    
    import java.io.IOException;
    
    @Component
    public class OrderConsumer {
    
        @Autowired
        private DispatcherService dispatcherService;
    
        @RabbitListener(queues = "order_dispatcher")
        public void messageReceive(String message) {
            JSONObject messageJSON = JSON.parseObject(message);
            Integer orderId = messageJSON.getInteger("orderId");
            dispatcherService.createDispatcher(orderId);
        }
    
    }
    
  2. 这里消费者一旦异常了,mq消息就会一直重发

    import com.alibaba.fastjson2.JSON;
    import com.alibaba.fastjson2.JSONObject;
    import com.zyy.rabbitmq.service.DispatcherService;
    import org.springframework.amqp.rabbit.annotation.RabbitListener;
    import org.springframework.beans.factory.annotation.Autowired;
    import org.springframework.stereotype.Component;
    
    import java.io.IOException;
    
    @Component
    public class OrderConsumer {
    
        @Autowired
        private DispatcherService dispatcherService;
    
        @RabbitListener(queues = "order_dispatcher")
        public void messageReceive(String message) {
            JSONObject messageJSON = JSON.parseObject(message);
            Integer orderId = messageJSON.getInteger("orderId");
            System.out.println(1 / 0); //这里异常,消息会一直重发
            dispatcherService.createDispatcher(orderId);
        }
    
    }
    
  3. 如何解决这个问题

    方法一:控制重发次数

    只要修改配置即可

      rabbitmq:
        username: admin
        password: admin
        virtual-host: /
        #    host: 39.108.49.252
        #    port: 5672
        addresses: 39.108.49.252:5672
        listener:
          simple:
            acknowledge-mode: manual # 开启手动ack,让程序去控制mq的消息的重发和删除和转移
            retry:
              enabled: true # 开启重试
              max-attempts: 10
              initial-interval: 2000ms
    

    方法二:try+catch+手动ack + 死信队列 + 人工干预 (推荐)

    配置队列

    import org.springframework.amqp.core.Queue;
    import org.springframework.context.annotation.Bean;
    import org.springframework.context.annotation.Configuration;
    
    import java.util.HashMap;
    import java.util.Map;
    
    @Configuration
    public class QueueConfig {
    
        @Bean("orderQueue")
        public Queue orderQueue() {
            //消息5秒后自动删除
            Map<String, Object> argMap = new HashMap<>();
            //死信队列
            argMap.put("x-dead-letter-exchange", "dead.order.fanout.exchange");
            return new Queue("order_dispatcher", true, false, false, argMap);
        }
    
        @Bean("deadQueue")
        public Queue deadQueue() {
            return new Queue("dead_order_dispatcher", true);
        }
    }
    

    配置交换机并绑定队列

    import org.springframework.amqp.core.Binding;
    import org.springframework.amqp.core.BindingBuilder;
    import org.springframework.amqp.core.FanoutExchange;
    import org.springframework.amqp.core.Queue;
    import org.springframework.beans.factory.annotation.Qualifier;
    import org.springframework.context.annotation.Bean;
    import org.springframework.context.annotation.Configuration;
    
    @Configuration
    public class FanoutRabbitConfig {
    
        @Bean("fanoutExchange")
        public FanoutExchange fanoutExchange() {
            return new FanoutExchange("order.fanout.exchange", true, false);
        }
    
        @Bean("deadExchange")
        public FanoutExchange deadFanoutExchange() {
            return new FanoutExchange("dead.order.fanout.exchange", true, false);
        }
    
    
        @Bean
        public Binding orderBinding(@Qualifier("orderQueue") Queue queue, @Qualifier("fanoutExchange") FanoutExchange exchange) {
            return BindingBuilder.bind(queue).to(exchange);
        }
    
        @Bean
        public Binding deadOrderBinding(@Qualifier("deadQueue") Queue queue, @Qualifier("deadExchange") FanoutExchange exchange) {
            return BindingBuilder.bind(queue).to(exchange);
        }
    }
    
    

    订单消费者

    import com.alibaba.fastjson2.JSON;
    import com.alibaba.fastjson2.JSONObject;
    import com.rabbitmq.client.Channel;
    import com.zyy.rabbitmq.service.DispatcherService;
    import org.springframework.amqp.rabbit.annotation.RabbitListener;
    import org.springframework.amqp.support.AmqpHeaders;
    import org.springframework.beans.factory.annotation.Autowired;
    import org.springframework.messaging.handler.annotation.Header;
    import org.springframework.stereotype.Component;
    
    import java.io.IOException;
    
    
    @Component
    public class OrderConsumer {
    
        @Autowired
        private DispatcherService dispatcherService;
    
        @RabbitListener(queues = "order_dispatcher")
        public void messageReceive(String message, Channel channel, @Header(AmqpHeaders.DELIVERY_TAG) long tag) {
            try {
                JSONObject messageJSON = JSON.parseObject(message);
                Integer orderId = messageJSON.getInteger("orderId");
                System.out.println(1 / 0); //这里异常,消息会一直重发
                dispatcherService.createDispatcher(orderId);
    
                channel.basicAck(tag, false);
    
                /**
                 * 解决这个问题
                 * 1.控制重发次数
                 * 2.try+catch+手动ack + 死信队列 + 人工干预
                 */
            } catch (Exception e) {
                /**
                 * 如果出现异常的情况,根据实际的情况去重发
                 * 重发一次后,丢失,还是存库根据自己的业务场景去决定
                 * long deliveryTag, 消息tag
                 * boolean multiple, 多条处理
                 * boolean requeue  重发
                 *   false 不会重发,会把消息打入死信队列
                 *   true 会死循环的重发,如果使用true 建议不用加try/catch 否则会造成死循环
                 */
                try {
                    channel.basicNack(tag, false, false);
                } catch (IOException ex) {
                    //监控
                    ex.printStackTrace();
                }
            }
        }
    
    }
    
    

    死信队列消费

    import com.alibaba.fastjson2.JSON;
    import com.alibaba.fastjson2.JSONObject;
    import com.rabbitmq.client.Channel;
    import com.zyy.rabbitmq.service.DispatcherService;
    import org.springframework.amqp.rabbit.annotation.RabbitListener;
    import org.springframework.amqp.support.AmqpHeaders;
    import org.springframework.beans.factory.annotation.Autowired;
    import org.springframework.messaging.handler.annotation.Header;
    import org.springframework.stereotype.Component;
    
    import java.io.IOException;
    
    @Component
    public class DeadOrderConsumer {
    
        @Autowired
        private DispatcherService dispatcherService;
    
        @RabbitListener(queues = "dead_order_dispatcher")
        public void messageReceive(String message, Channel channel, @Header(AmqpHeaders.DELIVERY_TAG) long tag) {
            try {
                JSONObject messageJSON = JSON.parseObject(message);
                Integer orderId = messageJSON.getInteger("orderId");
                dispatcherService.createDispatcher(orderId);
    
                channel.basicAck(tag, false);
            } catch (Exception e) {
                /**
                 * 死信队列还报错
                 * 那么:
                 * 人工干预
                 * 预警
                 * 同时把消息转移到别的存储db
                 */
                try {
                    channel.basicNack(tag, false, false);
                } catch (IOException ex) {
                    ex.printStackTrace();
                }
            }
        }
    
    }
    

总结

基于MQ的分布式事务解决方案优点:

  1. 通用性强
  2. 拓展方便
  3. 耦合度低,方案也比较成熟

基于MQ的分布式事务解决方案缺点:

  1. 基于消息中间件,只适合异步场景
  2. 消息会延迟处理,需要业务上能够容忍

建议

  1. 尽量去避免分布式事务
  2. 尽量将非核心业务做成异步

31、springboot整合rabbitMQ集群配置详解

 rabbitmq:
    addresses: 127.0.0.1:6605,127.0.0.1:6606,127.0.0.1:6705 #指定client连接到的server的地址,多个以逗号分隔(优先取addresses,然后再取host)
#    port:
    ##集群配置 addresses之间用逗号隔开
    # addresses: ip:port,ip:port
    password: admin
    username: 123456
    virtual-host: / # 连接到rabbitMQ的vhost
    requested-heartbeat: #指定心跳超时,单位秒,0为不指定;默认60s
    publisher-confirms: #是否启用 发布确认
    publisher-reurns: # 是否启用发布返回
    connection-timeout: #连接超时,单位毫秒,0表示无穷大,不超时
    cache:
      channel.size: # 缓存中保持的channel数量
      channel.checkout-timeout: # 当缓存数量被设置时,从缓存中获取一个channel的超时时间,单位毫秒;如果为0,则总是创建一个新channel
      connection.size: # 缓存的连接数,只有是CONNECTION模式时生效
      connection.mode: # 连接工厂缓存模式:CHANNEL 和 CONNECTION
    listener:
      simple.auto-startup: # 是否启动时自动启动容器
      simple.acknowledge-mode: # 表示消息确认方式,其有三种配置方式,分别是none、manual和auto;默认auto
      simple.concurrency: # 最小的消费者数量
      simple.max-concurrency: # 最大的消费者数量
      simple.prefetch: # 指定一个请求能处理多少个消息,如果有事务的话,必须大于等于transaction数量.
      simple.transaction-size: # 指定一个事务处理的消息数量,最好是小于等于prefetch的数量.
      simple.default-requeue-rejected: # 决定被拒绝的消息是否重新入队;默认是true(与参数acknowledge-mode有关系)
      simple.idle-event-interval: # 多少长时间发布空闲容器时间,单位毫秒
      simple.retry.enabled: # 监听重试是否可用
      simple.retry.max-attempts: # 最大重试次数
      simple.retry.initial-interval: # 第一次和第二次尝试发布或传递消息之间的间隔
      simple.retry.multiplier: # 应用于上一重试间隔的乘数
      simple.retry.max-interval: # 最大重试时间间隔
      simple.retry.stateless: # 重试是有状态or无状态
    template:
      mandatory: # 启用强制信息;默认false
      receive-timeout: # receive() 操作的超时时间
      reply-timeout: # sendAndReceive() 操作的超时时间
      retry.enabled: # 发送重试是否可用
      retry.max-attempts: # 最大重试次数
      retry.initial-interval: # 第一次和第二次尝试发布或传递消息之间的间隔
      retry.multiplier: # 应用于上一重试间隔的乘数
      retry.max-interval: #最大重试时间间隔

image-20230606230826164

Spring AMQP的主要对象

image-20230606230912149

32、rabbitMQ集群监控

参见:https://www.kuangstudy.com/zl/rabbitmq#1368199762003718146

33、rabbitMQ面试题分析

1、Rabbitmq 为什么需要信道,为什么不是TCP直接通信

1. TCP的创建和销毁,开销大,创建要三次握手,销毁要4次挥手
2. 如果不用信道,那应用程序就会TCP连接到Rabbit服务器,高峰时每秒成千上万连接就会造成资源的巨大浪费,而且底层操作系统每秒处理tcp连接数也是有限制的,必定造成性能瓶颈
3. 信道的原理是一条线程一条信道,多条线程多条信道同用一条TCP连接,一条TCP连接可以容纳无限的信道,即使每秒成千上万的请求也不会成为性能瓶颈

2、queue队列到底在消费者创建还是生产者创建?

1. 一般建议是在rabbitmq操作面板创建。这是一种稳妥的做法
2. 按照常理来说,确实应该消费者这边创建是最好,消息的消费是在这边。这样你承受一个后果,可能我生产在生产消息可能会丢失消息
3. 在生产者创建队列也是可以,这样稳妥的方法,消息是不会出现丢失
4. 如果你生产者和消费都创建的队列,谁先启动谁先创建,后面启动就覆盖前面的

文章评论