编程导航rabbitmq话题讨论

rabbitmq

17 参与
分享

快来分享你的内容吧~

点击登录,快来和大家讨论吧~
表情
图片
话题
打卡
综合
交流
文章
问答

05 Java后端面试八股合集 - 涵盖计网、JVM、Spring及RabbitMQ

# 00 引言 ## 写在最前面的话 本文是本人秋招时自用的八股准备文档,相较于八股网站繁多的内容,做了内容的压缩与高频考点的提炼,亲测覆盖面试 90% 以上的相关八股问题,分享出来以供大家参考。 另外本人已写 **Java后端完整版学习路线**(仔细讲解每个技术栈怎么学,学到什么程度可以投实习/面试)+ **各大厂真实面试问题和参考答案** + **Redis核心考点/面试真题深度梳理** + **MySQL高频考点深度梳理**(覆盖MySQL实习/校招面试95%以上的问题) 4篇高质量长文,如有需要欢迎大佬们自取(如果可以的话顺手点个赞👍),后续还会更新更多八股梳理、有趣的原创项目或实用的工具分享等内容,欢迎关注。 ### **其他优质内容导航:** **面试八股相关:** [01 非科班转码拿下大厂Offer,花费一天整理的Java后端完整版学习路线](https://www.codefather.cn/post/2005544965425446914) [02 盘点2025遇到的各大厂面试真题总结](https://www.codefather.cn/post/2005573237194469377) [03 一文吃透 Redis 核心考点,面试真题深度梳理](https://www.codefather.cn/post/2006203053207732225) [04 MySQL面试八股看这一篇就够了——深度梳理MySQL面试问题](https://www.codefather.cn/post/2008053544355131393) **工具分享**: [超实用的AI工具合集分享,赶快收藏起来吧 ](https://www.codefather.cn/post/2008402573266022402) **原创项目**: [01 基于用户画像和多模态需求驱动的多智能体推荐系统 - 技术文档](https://www.codefather.cn/post/2010000937266987009) (项目暂未开源,后续会开源,小伙伴可以关注一下) [02 ContentGuard Pro - 一套针对文本内容安全检测和风控的解决方案](https://www.codefather.cn/post/2010604558601961474) [Github链接](https://github.com/Mrchen-1600/Content-Guard-Pro) 03 TouchFish | 摸鱼神器(离线版)(最近即将推出,小伙伴们可以关注一下) 大致功能描述:是一款基于 Python 开发的高性能桌面隐私保护应用。它利用计算机视觉(CV)和离线语音识别(ASR)技术,实时监控用户周围环境。当检测到陌生人出现在摄像头内、用户离开品目前或触发特定语音关键词时,系统将毫秒级响应,执行静音、隐藏窗口并自动全屏打开用户提前设置的伪装工作软件/文件,为用户的“摸鱼”时光提供全方位保护。 ### 本文内容概述 **Notebook LM总结:** > 这份参考资料是一份详尽的计算机技术面试八股文合集,核心涵盖了计算机网络、JVM、Spring框架以及RabbitMQ四大技术模块。在网络层面,它详细拆解了从输入URL到页面展示的完整流程,并深入剖析了TCP可靠传输机制与HTTP/HTTPS的区别。针对Java虚拟机,文中深入探讨了内存区域划分、垃圾回收算法及对象创建过程,并提供了实用的JVM调优建议。框架部分重点解释了Spring IOC与AOP的核心原理,以及声明式事务在并发场景下的应用。最后,针对消息中间件,资料归纳了RabbitMQ如何处理消息丢失、重复消费及分布式事务一致性等实战难题。 # 计算机网络部分 01 从输入URL到页面展示到底发生了什么? ---------------------- 第一步是**用户输入URL**,并按下回车。此时,浏览器开始**解析URL**,将其分解为**协议**(例如HTTP或HTTPS)、**域名**、**路径**(例如/home)等部分。 第二步是进行**DNS解析**。浏览器通过DNS解析域名,**查找对应的IP地址**。若该域名的IP地址**已被缓存**,浏览器会**直接使用缓存的IP地址**;若**未缓存**,则浏览器会向DNS服务器**发送请求,获取IP地址**。 第三步是通过三次握手**建立TCP连接**。获取到IP地址后,浏览器通过**TCP/IP协议**与目标服务器**建立连接**。在HTTPS请求中,浏览器还会通过**SSL协议进行加密**,确保通信的安全性。 第四步是**客户端发送HTTP请求**。TCP连接建立后,浏览器会发送HTTP请求到服务器。请求包括**请求行、请求头和请求体**。**请求行包含请求方法(如GET、POST)、请求的路径和HTTP版本**。请求头包括目标域名、浏览器身份标识、可以处理的内容类型等信息。**请求体则用于发送数据(如POST表单提交数据)**。 第五步是**服务器处理请求**。服务器接收到请求后,首先会解析请求头并**根据请求的URL路径**选择合适的处理方式。 第六步是**服务器发送HTTP响应**。服务器处理完请求后,构造HTTP响应并返回给浏览器。**响应包含响应行、响应头和响应体**。**响应行指示处理结果(如200 OK、404 NotFound)**,响应头包含响应体的数据类型、响应体的长度等信息,**响应体则是实际的页面内容(HTML、CSS、JavaScript、图片等)**。 第七步是**浏览器解析和渲染页面**。浏览器接收到响应后,开始解析响应体中的HTML内容。 第八步是通过四次挥手**关闭TCP连接**。 02 TCP的粘包和拆包? ------------- ![](https://pic.code-nav.cn/post_picture/1734931576698986498/pCYodELIICOqH9Ky.webp) ![](https://pic.code-nav.cn/post_picture/1734931576698986498/osnCBERJKRi6esBU.webp) 03 常见的HTTP状态码? -------------- ![](https://pic.code-nav.cn/post_picture/1734931576698986498/LB09jvOnZeTDYsXX.webp) ![](https://pic.code-nav.cn/post_picture/1734931576698986498/B3LqJrMOA8fggsdI.webp) 04 TCP和UDP的区别? -------------- ![](https://pic.code-nav.cn/post_picture/1734931576698986498/meJ0dwU6ObTXfRyB.webp) 05 GET和POST的区别 -------------- GET和POST的5点主要区别: 第一个是数据传输方式的区别:GET请求**通过URL传递数据**,数据被附加在URL后面,以键值对的形式传输。而POST请求将数据**放在请求体中传递**,数据不会显示在URL中。 第二个是数据大小限制的区别:GET**请求的数据大小有限制**,通常为2至8KB。而POST请求的数据没有固定限制,**可以传输大量数据**,理论上只有Web服务器的配置限制。 第三个是安全性的区别:GET请求的数据通过URL传递,**数据泄露**给第三方,安全性较差。而POST请求的数据**保存在请求体中**,不会显示在URL中,**相对更安全**一些。 第四个是使用场景的区别:GET请求**通常用于获取资源或数据,适用于查询操作**。而POST请求通常用于**提交数据或进行数据更改**操作,适用于表单提交、用户注册、登录、文件上传等场景。 第五个是**幂等性**和缓存的区别:GET请求是幂等的,即同样的请求多次发送,服务器的响应不会发生变化,且可以被**缓存**。而POST请求是非幂等的,每次请求都会产生变化,因此不能缓存。 06 TCP超时重传机制为了解决什么问题? --------------------- ![](https://pic.code-nav.cn/post_picture/1734931576698986498/ExKduv8fDLZDkkmA.webp) 07 HTTP1.0、2.0和3.0的区别 --------------------- ![](https://pic.code-nav.cn/post_picture/1734931576698986498/TN38emSO7LwXvQco.webp) 08 服务器如何解析HTTP请求的数据 ------------------- ![](https://pic.code-nav.cn/post_picture/1734931576698986498/npRsrXkw0L8zuZDz.webp) 09 三次握手和四次挥手? ------------- ![](https://pic.code-nav.cn/post_picture/1734931576698986498/2xCB6Z8TmJgkE2dQ.webp) ![](https://pic.code-nav.cn/post_picture/1734931576698986498/uWF80TO7PTlZO3Wf.webp) 10 为什么要等2MSL?为什么是等待2MSL? ------------------------ (1)确保客户端最后的ACK能够被成功接收。因为如果这个ACK丟失了,服务器没有收到确认包,会重新发送FIN报文,而MSL是TCP报文在网络中可以存活的最大时间,服务器重发FIN,客户端收到之后重发ACK,这一来一回就需要2MSL的时间。 (2)防止旧的报文干扰新的连接。TCP连接关闭之后,可能会有一些延迟的或者已经失效的报文还在网络中传输,如果我们立即用相同的IP地址和端口建立新的连接,可能会受到这些旧的报文的干扰。 11 TCP实现可靠传输的原理 --------------- 首先是**连接管理机制**,TCP通过**三次握手建立连接**,确保双方通信正常,并用**四次挥手终止连接**,防止数据残留或资源浪费。 第二个是数据分块与序号标识机制,发送端将**数据分割为合适大小的报文段**,每个段分配**唯一序号**,**标识数据的顺序**;接收端**通过序号重组乱序到达的段**,确保数据完整性。 第三个是**确认应答与超时重传机制**,接收端对每个接收到的段返回确认应答,发送端如果超时未收到ACK就会触发超时重传,解决数据丢失问题。 第四个是**流量控制**机制,接收端**通过窗口大小告知发送端可接收的数据量**,**避免缓冲区溢出**。滑动窗口机制允许连续发送多个段,提升传输效率。 第五个是**拥塞控制**机制,能够根据网络负载情况**动态调整发送速率**,防止网络瘫痪。 12 HTTP怎么实现流量控制(滑动窗口算法) ----------------------- ![](https://pic.code-nav.cn/post_picture/1734931576698986498/ziNUt9vkTcw7PW5J.webp) ## 13 TCP每次连接时序列号都一样吗?为什么不一样,有什么作用? ![](https://pic.code-nav.cn/post_picture/1734931576698986498/hudQe8PQ1V4ahl0M.webp) ![](https://pic.code-nav.cn/post_picture/1734931576698986498/HslMLHido6eHVqJW.webp) ![](https://pic.code-nav.cn/post_picture/1734931576698986498/5aj5cM7x4RgGMDme.webp) 14 HTTPS协议和HTTP协议的区别? --------------------- (1)数据传输安全性: http:明文传输,容易被窃听、篡改 https:通过SSL/TSL协议对数据进行加密传输,提供数据机密性和完整性保障。 (2)端口号: http:默认端口号80 https:默认端口号443 (3)性能: http:无加密过程,连接建立速度稍快。 https:基于http上又加了SSL或TSL协议来实现的加密传输,加解密过程增加了计算开销,握手时间较长。 ![](https://pic.code-nav.cn/post_picture/1734931576698986498/xOMUGyCHCGhpm1rC.webp) 15 DDOS攻击 --------- DDOS攻击(Distributed Denial of Service,**分布式拒绝服务攻击**)是一种通过大量恶意流量淹没目标服务器、网络或服务,使其无法正常响应合法用户请求的网络攻击方式。 **基本原理** **拒绝服务(DoS)**:攻击者通过耗尽目标的带宽、计算资源(如CPU、内存)或应用处理能力,导致服务瘫痪。 **分布式(Distributed)**:攻击流量来自全球大量被控制的设备(如僵尸网络中的电脑、IoT设备等),而非单一来源,难以追踪和防御。 **常见攻击类型** **流量洪泛** 例如:UDP洪水,通过垃圾流量塞满目标带宽。 **协议攻击** 例如:SYN洪水(耗尽TCP连接资源)、DNS放大攻击(利用DNS协议缺陷放大流量)。 **应用层攻击** 例如:HTTP洪水(模拟大量合法请求耗尽服务器资源),更隐蔽且难以识别。 **防御措施** **流量清洗**:通过云安全服务(如Cloudflare、阿里云高防IP)过滤恶意流量。 **黑名单/IP限速**:识别并拦截异常IP。 **冗余架构**:分布式服务器分散流量压力。 # JVM部分 01 GC 怎么知道哪些是垃圾?(垃圾搜索的算法) ------------------------- 垃圾回收需要知道哪些对象是垃圾(垃圾的搜集算法),主要是两种方式: (1)**引用计数法:**每个对象都有一个**引用计数器**,每当有一个引用指向他,计数器就加1,所以GC只要去看对象的计数器是不是0就行,是0的就可以直接回收,但是这个方法最大的问题就是,如果**两个对象相互引用**,那就永远不能被回收,所以就需要方法2 (2)**可达性分析算法** 从**GC Roots**出发,看看哪些对象是可达的,可达的即存在引用,不可达的就可以直接回收(GC Roots可以是栈中引用的对象,类静态属性引用的对象,常量引用的对象以及本地方法引用的对象) 02 GC Roots包含哪些对象 ----------------- ![](https://pic.code-nav.cn/post_picture/1734931576698986498/pLO04pQRrcUFNz2g.webp) ![](https://pic.code-nav.cn/post_picture/1734931576698986498/cSFMxY6ffjFb6ALS.webp) 标记出垃圾之后就是进行回收了,下面是回收的算法: 03 GC垃圾回收算法 ----------- **1.标记清除算法** 当垃圾回收器进行内存扫描后会标记出所有的垃圾,然后**直接清除这些带标记**的垃圾对象即可,但是因为垃圾在内存上是分散的,这样清理就会**产生大量的内存碎片**,使得内存利用率越来越低; **2.复制算法** **准备两块一模一样的内存空间**,当第一块剩余空间不足时,可以将**所有需要保留的对象拷贝**至另一块空内存,然后将**前一块内存全部清空**,这样即做到了垃圾回收,又做到了碎片整理,但是这样**内存空间就会有一半被浪费**。 **3.标记整理法** 标记整理算法就是在清理垃圾的基础上,多了一步碎片整理的工作,因为整理是比较耗时的,所以显然这种垃圾回收机制不适合高频率的执行。 知道了常见的垃圾回收算法,再介绍下常见的垃圾回收器: 04 JVM常见的垃圾回收器 -------------- 最早期的就是**Serial和SerialOld**,但是这种垃圾回收器是**单线程的**,不支持并发,**开始垃圾回收后所有用户线程都必须暂停**。 后面是**parallel**的垃圾回收器,垃圾回收可以**多线程执行**,但是当开始垃圾回收的时候依然会**触发所有用户线程暂停**。 再之后就是**CMS**,CMS虽然在**初始标记的时候也会触发用户线程暂停**,但是因为CMS**初始标记只会标记第一层的根对象**,所以时间很短,**真正的标记阶段,打标记是和用户线程并行的**;为了避免并行过程中可能的错标,还有第三个**重新标记的阶段**,**这个阶段也会让用户线程暂停,然后通过多线程的方式去检查并且修正错误标记。**但是这个算法也有问题,就是**清理垃圾的时候用户线程也在运行,因此新产生的垃圾没办法及时清理**,需要等下一次回收。 再就是**G1回收器**,G1回收器会**把整个堆内存划分成若干相等大小的区域**,然后**对这些区域进行价值排序**,垃圾越多,回收需要的时间也越多,但是G1回收器可以**根据我们设定的用户线程暂停时间来调整策略**,会尽量**满足我们设定的时间去回收价值高的区域**。另外**G1还把大对象单独存放在了一个区域**,避免了整理时候**频繁移动大对象**。 补充一个新生代和老年代的知识: 05 新生代和老年代? ----------- 区分新生代和老年代主要是为了提高垃圾回收效率,大多数的对象存活时间短,很快就会变成垃圾不再使用,这些短生命周期的对象就会分配在新生代;少部分对象长期存活,不会很快被回收,就晋升到老年代。 针对不同的分区的垃圾回收算法也不一样,**新生代通常采用复制算法**,**老年代**通常采用**标记整理算法或标记清除算法**。 06 什么是FullGC?什么情况下触发?怎么解决? -------------------------- Full GC(完全垃圾回收)是指对整个JVM**堆(包括新生代、老年代)**以及**方法区(元空间)**进行的全面垃圾回收。Full GC会暂停所有应用线程(Stop-The-World),通常耗时更长。 **触发的场景**: 1. 老年代空间不足,且通过old GC仍然不足 2. 元空间或者永久代内存不足 3. 使用了System.gc()命令 4. 新生代对象要晋升到老年代,但是老年代空间不够 **频繁full GC的问题**: 长时间的Stop-The-World暂停会导致应用响应变慢,大量的CPU时间用于GC而非业务处理,也会导致业务吞吐量下降,可能会导致OOM或者服务超时 **解决**: 增加堆大小、调整新生代与老年代的比例、选择合适的垃圾回收器(比如G1替代CMS)、代码优化(及时释放不再使用的对象,减少大对象的频繁创建) 07 JVM的内存区域? ------------ 首先是**线程共享**的部分,一共有两个: 一个是**堆(Heap)**,所有**对象实例和数组**都在这里分配内存,垃圾回收器(**GC**)也主要在堆中工作。堆中还包含了**字符串常量池**(String Constant Pool)。 另一个是**元空间(Method Area)**:用于**存储类信息、常量、静态变量、方法字节码**等。其中运行时常量池(Runtime Constant Pool)是元空间的一部分,用于存储编译期生成的各种字面量和符号引用。 ### 字符串常量池的作用? **为什么需要字符串常量池?** String s1 = "Hello"; String s2 = "Hello"; String s3 = new String("Hello"); 如果没有常量池:s1会创建一个新的String对象。s2会再创建一个内容完全相同的新String对象。s3通过new关键字,毫无疑问也会创建一个新对象。这样,内存中就会有三个内容完全相同的"Hello"对象,这是极大的浪费。 字符串常量池就是为了解决这个问题而生的。**工作原理(核心机制):**当在代码中直接使用字符串字面量(用双引号包裹)时,例如String s = "Hello";,JVM会首先去字符串常量池中查找是否已经存在一个内容为"Hello"的字符串对象。 **如果存在**:JVM不会创建新的对象,而是直接返回池中已有对象的引用。这样,所有相同的字面量都指向同一个内存地址。 **如果不存在**:JVM会在字符串常量池中创建一个新的String对象,内容为"Hello",然后返回这个新对象的引用。 new String()**的创建**:当使用new关键字(如String s = new String("Hello");)时,JVM的行为会有所不同:new关键字会**强制**在Java堆的**非常量池区域**创建一个全新的、独立的String对象。这个新对象的内容会初始化为"Hello",但它和常量池中的那个"Hello"是两个不同的对象。 ![](https://pic.code-nav.cn/post_picture/1734931576698986498/lkFjjQ5ioMyI7U3c.webp) 然后是**线程私有**的部分,一共有三个, 第一个是**虚拟机栈**(VM Stack),**每个线程启动时都会创建一个虚拟机栈**,它存储方法调用过程中产生的栈帧,包括**局部变量、方法返回地址**等,**每个方法调用都会创建一个新的栈帧**,方法执行结束后栈帧出栈。 第二个是**本地方法栈**(Native Method Stack),专门用于存储**本地方法**(Native Method)的调用信息,与虚拟机栈类似,但用于JNI(Java Native Interface)调用。 第三个是**程序计数器**(Program Counter Register),**记录当前线程正在执行的字节码指令地址**。它是JVM运行时最小的内存区域,每个线程都有一个独立的程序计数器。 在JDK1.8时 JVM的内存结构主要有两点不同: 一个是方法区(Method Area)在JDK 1.8被替换为元空间(Metaspace)实现,且元空间使用本地内存(原因:元空间可以动态调整大小,能够避免Out of Memory的错误,并且提高了GC回收效率) 另一个是运行时常量池(Runtime Constant Pool)在JDK1.7属于方法区的一部分,而在JDK 1.8变成元空间的一部分。 08 对象创建的过程了解吗? -------------- 第一步是进行**类加载检查**,当程序执行到new指令时,JVM会**先检查对应的类是否已经被加载**、解析和初始化过。如果类尚未加载,JVM会按照类加载机制(加载、验证、准备、解析、初始化)完成类的加载过程。这一步**确保了类的元信息(如字段、方法等)已经准备好**,为后续的对象创建奠定基础。 第二步是进行**内存的分配**,JVM会为新对象分配内存空间。对象所需的内存大小在类加载完成后就可以确定,因此**分配内存的过程就是从堆中划分一块连续的空间**,主要有两种方式: 一种是通过指针碰撞,**如果堆中的内存是规整的**(已使用和空闲区域之间有明确分界),JVM可以**通过移动指针来分配内存**。另一种是通过空闲列表,如果**堆中的内存是碎片化的**,JVM会维护一个**空闲列表**,记录可用的内存块,并从中分配合适的区域。 第三步是将**零值初始化**,JVM会对分配的内存空间进行初始化,将其所有字段设置为零值(如int为0,boolean为false,引用类型为null)。这一步确保了对象的实例字段在未显式赋值前有一个**默认值**,从而避免未初始化的变量被访问。 第四步是**设置对象头**,其中包含Mark Word、Klass Pointer和数组长度。Mark Word用于**存储对象的哈希码、GC分代年龄**等信息。Klass Pointer指向对象所属类的元数据(即Person.class的地址)。 第五步是**执行构造方法**,用**<init>方法完成对象的初始化**。构造方法会根据代码逻辑对对象的字段进行赋值,并调用父类的构造方法完成继承链的初始化。这一步完成后,对象才真正可用。 09 JVM相关的配置与调优专题 ---------------- ### **一、堆配置:** **1.1.配置项** `-Xms`:初始堆大小 `-Xmx`:最大堆大小 `-XX:NewSize=n`:设置年轻代大小 `-XX:NewRatio=n`:设年轻代和年老代的比值。如:为3表示年轻代和年老代比值为1:3,年轻代占整个年轻代年老代和的1/4,默认为2 `-XX:SurvivorRatio=n`:年轻代中Eden区与两个survivor区的比值,注意Survivor区有两个,默认8。如:3表示Eden:3 Survivor:2,一个Survivor区占整个年轻代的1/5 `-XX:MaxPermSize=n`:设置持久代大小 **1.2.说明** 1、一般初始堆和最大堆设置一样,因为:现在内存不是什么稀缺的资源,但是如果不一样,从初始堆到最大堆的过程会有一定的性能开销,所以一般设置为初始堆和最大堆一样。64位系统理论上可以设置为无限大,但是一般设置为4G,因为如果再大,JVM进行垃圾回收出现的暂停时间会比较长,这样全GC过长,影响JVM对外提供服务,所以不能太大。一般设置为4G。 2、`-XX:NewRaio`和`-XX:SurvivorRatio`这两个参数,第一是设置年轻代的大小,二个是设置年轻代的比值、理论上设置一个即可以满足需求(因为有默认值) **1.3. 概念解释** 年轻代:包括Eden区和Survivor区,用于管理新创建的对象。 老年代:用于存放从年轻代中存活下来的对象。 * Eden区:新对象首先被分配到这里。 * Survivor区:用于存放从Eden区中存活下来的对象,通过两个Survivor区的交替使用减少内存碎片。 **年轻代(Young Generation)** 年轻代是对象最初被分配的地方。大多数对象在年轻代中创建,并且大多数对象在年轻代中死亡(即不再被引用)。年轻代的主要特点是: Eden区:新创建的对象首先被分配到Eden区。Eden区是年轻代的主要部分,通常占据年轻代的大部分空间。 Survivor区:Survivor区分为两个部分,通常称为From Survivor区和To Survivor区。在每次Minor GC(年轻代垃圾回收)后,存活的对象会被移动到Survivor区中的一个,而另一个Survivor区则被清空。这种设计有助于减少内存碎片。(复制算法) **老年代(old Generation)** 老年代用于存放从年轻代中存活下来的对象。当对象在年轻代中经过多次Minor GC后仍然存活,它们会被晋升到老年代。老年代的特点是: Major GC(Full GC):老年代的垃圾回收称为Major GC或Full GC。Full GC通常比Minor GC更耗时,因为它需要扫描整个堆内存。 **Eden区与Survivor区** Eden区:新对象首先被分配到Eden区。当Eden区满时,会触发一次Minor GC,将存活的对象移动到Survivor区中的一个(通常是From Survivor区)。 Survivor区:Survivor区用于存放从Eden区中存活下来的对象。在每次Minor GC后,存活的对象会被移动到另一个Survivor区(To Survivor区),而原来的Survivor区(From Survivor区)则被清空。这种设计有助于减少内存碎片,提高内存利用率。 ### **二、调优总结** **年轻代大小选择:** * **响应时间优先的应用:** 尽可能设置大,直到接近系统的最低响应时间限制(根据实际情况选择)。在此种情况下,年轻代收集发生的频率也是最小的。同时减少到达年老代的对象。 * **吞吐量优先的应用:** 尽可能的设置大,可能到达Gbit的程度,因为对响应时间没有要求,垃圾收集可以并行进行,一般适合8核CPU以上应用。 **年老代大小选择:** * **响应时间优先的应用:** 年老代使用并发收集器,所以其大小需要小心设置,一般要考虑并发会话率和会话持续时间等一些参数。如果堆设置小了,可能会造成内存碎片、高回收频率以及应用暂停而使用传统的标记清除方式;如果堆大了,则需要较长的收集时间。最优化的方案,一般需要参考以下数据获得: * 1、并发垃圾收集信息 * 2、持久代并发收集次数 * 3、传统GC信息 * 4、花在年轻代和年老代回收上的时间比例,减少年轻代和年老代花费的时间,一般会提高应用的效率。 * **吞吐量优先的应用:** 一般吞吐量优先的应用都有一个很大的年轻代和一个较小的年老代。原因是,这样可以尽可能回收掉大部分短期对象,减少中期对象,而年老代尽存放长期存活的对象 **较小堆引起的碎片问题:** 因为年老代的并发收集器使用标记、清除算法,所以不会对堆进行压缩。当收集器回收时,他会把相邻的空间进行合并,这样可以分配给较大的对象。但是当堆空间较小时,运行一段时间以后,就会出现“碎片”,如果并发收集器找不到足够的空间,那么并发收集器将会停止,然后使用传统的标记、清除方式进行回收。如果出现“碎片”,可能需要进行如下配置: `-XX:+UseCMSCompactAtFullCollection`:使用并发收集器时,开启对年老代的压缩 `-XX:CMSFullGCsBeforeCompaction=0`:上面配置开启的情况下,这里设置多少次FullGc后,对年老代进行压缩。 ### **三、内存泄露检查** 根据垃圾回收前后情况对比,同时根据对象引用情况(常见的集合对象引用)分析,基本都可以找到泄漏点。 **持久代占满处理:** 1、`-XX:MaxPermSize=16m`设置持久代大小 2、换JDK、比如:JRocket **系统内存被占满:** 一般是因为没有足够的资源产生线程造成的,系统创建线程时,除了要在Java堆中分配内存外,操作系统本身也需要分配资源来创建线程。因此,当线程数量大的一定程度以后,堆中或许还有空间,但是操作系统分配不出资源来了,出现异常; 分配给Java虚拟机的内存越多,系统剩余的资源就越少,因此,当系统内存固定时,分配给Java虚拟机的内存越多,那么,系统总共能够产生的线程也就越少,两者成反比。同时,可以通过修改-Xss来减少分配给单个线程的空间,也可以增加系统总共生产的线程数。 ### **四、GC使用与性能优化管理** (1)不要显式调用System.gc()。此函数建议JVM进行主动GC,虽然只是建议而非一定,但很多情况下它会触发主GC,从而增加主动GC的频率、也即增加了间歇性停顿的次数。大大的影响系统性能。 (2)尽量减少临时对象的使用。临时对象在跳出函数调用后,会成为垃圾,少用临时变量就相当于减少了垃圾的产生,从而减少了主GC的机会。 (3)对象不用时最好显式置为Null。一般而言,为Null的对象都会被作为垃圾处理,所以将不用的对象显式地设为Null,有利于GC收集器判定垃圾,从而提高了GC的效率。 (4)尽量使用StringBuffer,而不用String来累加字符串。由于String是固定长的字符串对象,累加String对象时,并非在一个String对象中扩增,而是重新创建新的String对象,如`Str5=Str1+Str2+Str3+Str4`,这条语句执行过程中会产生多个垃圾对象,因为对次作“+”操作时都必须创建新的String对象,但这些过渡对象对系统来说是没有实际意义的,只会增加更多的垃圾。避免这种情况可以改用StringBuffer来累加字符串,因StringBuffer是可变长的、它在原有基础上进行扩增,不会产生中间对象。 (5)能用基本类型如int,long,就不用Integer,Long对象。基本类型变量占用的内存资源比相应对象占用的少得多,如果没有必要,最好使用基本变量。 (6)尽量少用静态对象变量。静态变量属于全局变量,不会被GC回收,它们会一直占用内存 (7)注意分散对象创建或删除的时间,集中在短时间内大量创建新对象,特别是大对象,会导致突然需要大量内存、JVM在面临这种情况时,只能进行主GC,以回收内存或整合内存碎片,从而增加主GC的频率。集中删除对象,道理也是一样的。 ### **五、JVM调优参数参考** 1.针对JVM堆的设置,一般可以通过`-Xms -Xmx`限定其最小、最大值,为了防止垃圾收集器在最小、最大之间收缩堆而产生额外的时间,通常把最大、最小设置为相同的值; 2.年轻代和年老代将根据默认的比例(1:2)分配堆内存,可以通过调整二者之间的比率NewRadio来调整二者之间的大小,也可以针对回收代 比如年轻代,通过: `-XX:newSize-XX:MaxNewSize`来设置其绝对大小。同样,为了防止年轻代的堆收缩,我们通常会把`-XX:newSize-XX:MaxNewSize`设置为同样大小。 3.年轻代和年老代设置多大才算合理 * 更大的年轻代必然导致更小的年老代,大的年轻代会延长普通GC的周期,但会增加每次GC的时间;小的年老代会导致更频繁的FullGC * 更小的年轻代必然导致更大年老代,小的年轻代会导致普通GC很频繁,但每次的GC时间会更短:大的年老代会减少FullGC的频率 如何选择应该依赖应用程序对象生命周期的分布情况:如果应用存在大量的临时对象,应该选择更大的年轻代;如果存在相对较多的持久对象,年老代应该适当增大。但很多应用都没有这样明显的特性。 在抉择时应该根据以下两点: (1)本着FullGC尽量少的原则,让年老代尽量缓存常用对象,JVM的默认比例1:2也是这个道理。 (2)通过观察应用一段时间,看其他在峰值时年老代会占多少内存,在不影响FullGC的前提下,根据实际情况加大年轻代,比如可以把比例控制在1:1。但应该给年老代至少预留1/3的增长空间。 4.在配置较好的机器上(比如多核、大内存),可以为年老代选择并行收集算法-`XX+UseParallelOldGC`。 5.线程堆栈的设置:每个线程默认会开启1M的堆栈,用于存放栈帧、调用参数、局部变量等,对大多数应用而言这个默认值太了,一般256K就足用。 理论上,在内存不变的情况下,减少每个线程的堆栈,可以产生更多的线程,但这实际上还受限于操作系统。 ### **六、触发Full GC的几种情况** **1.老年代空间不足** (1)年轻代晋升:当年轻代中的对象经过多次Minor GC后仍然存活,会被晋升到老年代。如果老年代空间不足,无法容纳这些晋升的对象,就会触发Full GC。 (2)大对象直接分配:如果应用程序直接分配大对象(超过年轻代Eden区的大小),这些对象会直接进入老年代。如果老年代空间不足,也会触发Full GC。 **2.方法区(元空间或永久代)空间不足** (1)类加载:如果应用程序动态加载大量类,导致方法区(Metaspace或永久代)空间不足,会触发Full GC。 (2)元数据回收:Metaspace或永久代中的元数据(如类信息、方法信息等)不再被引用时,需要通过Full GC进行回收。 **3.System.gc()调用** 显式调用:应用程序中显式调用System.gc()或Runtime.getRuntime().gc()会建议JVM执行Full GC。不过JVM不一定会立即执行,具体行为取决于JVM的实现和配置。 **4.垃圾回收器策略** (1)CMS GC的并发模式失败:在使用CMS(Concurrent Mark-Sweep)垃圾回收器时,如果在并发标记过程中老年代空间不足,会触发Full GC。 (2)G1 GC的疏散失败:在使用G1(Garbage-First)垃圾回收器时,如果在疏散(Evacuation)阶段无法找到足够的空闲区域来存放存活对象,会触发Full GC。 **5.堆内存分配失败** 内存泄漏:如果应用程序存在内存泄漏,导致堆内存中大量对象无法被回收,最终堆内存耗尽,会触发Full GC # Java集合部分 01 HashMap的原理 ------------- HashMap在jdk1.7和1.8实现上是有些不一样的,先介绍1.7。在1.7中,底层是通过**数组+链表**来实现的,当我们插入元素的时候,会计算key的hash值,也就是**对数组长度取模**,得到插入位置的下标,如果该位置为空,就插入元素,**如果不为空,就会以链表的形式存储**,会去遍历链表,如果找到了相同的key,就替换value,表示修改操作;如果没有相同的key,就把新的entry**插入到链表的头部**。在1.7中,扩容机制是,比如默认数组容量是16,加载因子0.75,所以当我们插入第13个元素的时候就会触发扩容,**扩容会重新对所有元素进行hash计算,去把元素放到新的数组中去**。 但是1.7中存在一些问题:首先就是**哈希冲突比较严重的时候链表会变得很长**,链表的查询效率是O(n),这就会影响性能。另外,**头插法虽然比较快,但是在多线程环境可能就会形成环形链表,陷入死循环**(比如现在链表是A->B->C,线程A和线程B同时去操作链表,如果线程A先去操作时候发现需要扩容,通过头插法扩容A先放入新数组,然后是B和C,顺序就变成了C->B->A,这个时候线程B再去操作,线程B还以为是原链表,即A指向B,但是现在实际已经变成了B指向A,就形成了死循环);再就是扩容,1.7是对所有元素重新计算,这个也比较复杂。 针对这3个问题,首先一点,1**.8改成了尾插法**,**扩容不会反转链表**,所以避免了死循环的产生;第二点,**1.8中采用数组+链表+红黑树的结构**。当我们插入数据,如果对应位置已经有元素,会先存储成链表,但是如果**链表的长度已经等于8**了,就需要看**数组长度是否大于64**,**如果没有就优先扩容数组**;**如果数组元素已经大于64,就把链表转换成红黑树存储**,这个转换的契机就是链表长度大于8,然后**如果元素少于6个,就从二叉树再退化回链表**。1.8里的扩容契机和1.7一样,除了前面提到的链表长度那里以外,也是判断数组存储的元素是否超过临界值。但是1.8的扩容不是全部重新计算hash,而是**通过元素的hash和老数组的长度进行&运算来计算出元素是处于高位还是低位**,**如果结果是0就把元素留在原来位置不移动,否则就移动到原来索引位置加上老数组容量的位置去**,显著提升了扩容的速度。 ![](https://pic.code-nav.cn/post_picture/1734931576698986498/BQs0w7d9uF41OrSl.webp) 注意:声明初始容量的时候需要考虑集合的实现类型,如果是hashmap或者hashset(实际就是hashmap,只不过value为空),那初始容量实际要声明成100W/扩容因子(默认0.75),如果是非hashmap实现,比如arraylist,直接声明成100W就可以了,因为arraylist是存满再扩容,没有扩容因子。 02 ArrayList和LinkedList的区别? --------------------------- ![](https://pic.code-nav.cn/post_picture/1734931576698986498/Dn5EXGcLxbwBrVWW.webp) 03 常见的集合有哪些? ------------ ![](https://pic.code-nav.cn/post_picture/1734931576698986498/YLiRbWx5BxIGcULB.webp) 04 1亿量级的ArrayList数据去重? ---------------------- (1)直接hashset,但是注意初始化容量,避免频繁扩容(可以初始化成ArrayList.size() / 0.75 + 1) (2)排序+遍历去重,排序完之后,一遍扫过去,只要和前一个不一样就保留; (3)bitmap,先遍历一遍arraylist,用一个二进制位去把arraylist里元素值对应位置的bit值改为1,遍历完再遍历位图,收集所有标记为1的位置下标(如果是0-1亿之间的数字,bitmap只需要12MB的内存,如果是int范围512MB也够了) # SSM框架部分 01 Spring IOC ------------- IoC即控制反转。 例如:现有类A依赖于类B 传统的开发方式:往往是在类A中手动通过new关键字来new一个B的对象出来;使用IoC思想的开发方式:不通过new关键字来创建对象,而是通过IoC容器(Spring框架)来帮助我们实例化对象。我们需要哪个对象,直接从IoC容器里面去取即可。 从以上两种开发方式的对比来看:我们“丧失了一个权力” (创建、管理对象的权力),从而也得到了一个好处(不用再考虑对象的创建、管理等一系列的事情) 为什么叫控制反转? 控制:指的是对象创建(实例化、管理)的权力 反转:控制权交给外部环境(IoC容器) ### IoC解决了什么问题? IoC的思想就是对象之间不互相依赖,由第三方容器来管理相关资源。这样有什么好处呢? 1、降低了对象之间的耦合度或者说依赖程度; 2、资源变的容易管理;比如用Spring容器提供的话很容易就可以实现一个单例。 例如:现有一个针对User的操作,利用Service和Dao两层结构进行开发,在没有使用IoC思想的情况下,Service层想要使用Dao层的具体实现的话,需要通过new关键字在UserService的实现类中手动new出UserDao的具体实现类。如果后续接到新的需求,针对UserDao接口需要开发另一个新的实现类,我们就需要手动修改UserService实现类中new的对象。如果有许许多多的地方都引用了UserDao的具体实现的话,那修改起来就会非常的繁琐。 使用IoC的思想,我们将对象的控制权(创建、管理)交由IoC容器去管理,我们在使用的时候直接向IoC容器 “要”就可以了(通过注解去标识,按照类型/名称去进行注入)。 ### IoC和DI有区别么? IoC是一种设计思想或者说是某种模式。这个设计思想就是**将原本在程序中手动创建对象的控制权交给第三方比如IoC容器。**对于我们常用的Spring框架来说,IoC容器实际上就是个Map(key,value),Map中存放的是各种对象。 IoC最常见的实现方式叫做依赖注入简称DI。 02 Spring自动装配原理 --------------- ![](https://pic.code-nav.cn/post_picture/1734931576698986498/ALVvh0rZcAqWrQDT.webp) 通过注解或者一些简单的配置就能在Spring Boot的帮助下快速实现某块功能。 在传统Spring中,我们需要在XML或Java配置中显式地定义很多Bean(如数据源等)。而在Spring Boot中,只引入了特定的依赖,Spring Boot就会自动配置好这些组件。举个例子:要使用JDBC,只需在pom.xml中引入spring-boot-starter-jdbc依赖,并在yml文件中配置一些连接信息即可。 **自动装配的实现:** 核心机制:条件化装配 这是自动装配的基石。Spring Boot不会盲目地配置所有东西,它只在某些**条件满足**的情况下才进行配置。这是通过一系列@Conditional注解来实现的。 条件注解: 1. 如果没有引入web启动依赖,所有与web相关的自动配置类都不会被加载; 2. 用户自定义的加了@Bean注解的都会被优先加载; 3. 在yml等配置文件中我们可以自定义配置属性,比如内置的tomcat启动端口是8080,如果我们配置成其他端口,就会按照我们的配置进行自动装配。 自动装配的过程可以概括为以下几个步骤: 1. **启动注解:**@SpringBootApplication 每个Spring Boot主类上都标有@SpringBootApplication,它是一个复合注解,其中@EnableAutoConfiguration注解就启用自动配置的。 1. **启用自动配置:**@EnableAutoConfiguration 这个注解的作用是启用Spring Boot的自动配置机制。它背后是通过一个import选择器的类实现的(AutoConfigurationImportSelector.class)。 1. **加载自动配置列表:** 会通过import选择器类读Classpath下所有JAR包中的META-INF/spring.factories文件。 1. **关键文件:** META-INF/spring.factories 这个文件是一个标准的Java配置文件,内容是**键=值**对的形式。这个列表定义了**所有可能被自动配置的类**。 1. **过滤与条件判断:** Import选择器不会直接加载所有可能的配置类。而是通过**条件注解**对列表中的配置类进行筛选。只有满足条件的配置类,才会被真正解析,将其中的Bean定义加载到Spring容器中。 03 Spring AOP ------------- AOP(Aspect Oriented Programming)即面向切面编程,AOP是OOP(面向对象编程)的一种延续,二者互补,并不对立。 AOP的目的是将一些分散在多个类或对象中的公共行为(如日志记录、事务管理、权限控制、接口限流等)从核心业务逻辑中分离出来,通过动态代理技术,实现代码的复用和解耦,提高代码的可维护性和可扩展性。 AOP之所以叫面向切面编程,是因为它的核心思想就是将横切关注点从核心业务逻辑中分离出来,形成一个个的切面(Aspect)。 ### AOP常见的通知(增强)类型 ![](https://pic.code-nav.cn/post_picture/1734931576698986498/YVfKWVMQS1yHXpFm.webp) ### AOP解决了什么问题? OOP不能很好地处理一些分散在多个类或对象中的公共行为(如日志记录、事务管理、权限控制、接口限流、接口幂等等),这些行为通常被称为横切关注点。如果我们在每个类或对象中都重复实现这些行为,那么会导致代码的冗余、复杂和难以维护。 AOP可以将横切关注点(如日志记录、事务管理、权限控制、接口限流、接口幂等等)从核心业务逻辑中分离出来,实现关注点的分离。 比如日志记录,没有AOP之前,我们需要对需要加日志记录的地方挨个去写代码,全是重复的逻辑;但是有AOP技术之后,我们就可以把日志记录的逻辑封装成一个切面,然后通过切点和通知来指定具体在哪些方法中加入日志。在指定方法中只需要加一行注解就可以实现日志记录。 ### AOP的应用场景 **日志记录:** 自定义日志记录注解,利用AOP,给业务方法上加一行注解即可实现日志记录。 **性能统计:** 利用AOP在目标方法的执行前后统计方法的执行时间,方便优化和分析。 **事务管理:** @Transactional注解可以让Spring为我们进行事务管理比如回滚异常操作,免去了重复的事务管理逻辑。@Transactional注解就是基于AOP实现的。 **权限控制:** 利用AOP在目标方法执行前判断用户是否具备所需要的权限,如果具备,就执行目标方法,否则就不执行。 **接口限流:** 利用AOP在目标方法执行前通过具体的限流算法对请求进行限流处理。 ### AOP动态代理 Spring AOP是基于动态代理的,如果要代理的对象,实现了某个接口,那么Spring AOP会使用**JDK代理**,去创建代理对象,而对于没有实现接口的对象,就无法使用JDK去进行代理了,这时候Spring AOP会使用CGLIB生成一个被代理对象的子类来作为代理。 04 Spring事务 ----------- 事务是逻辑上的一组操作,要么都执行,要么都不执行。 我们系统的每个业务方法可能包括了多个原子性的数据库操作,比如转账就是经典的事务场景。事务能否生效数据库引擎是否支持事务是关键。比如常用的MySQL数据库默认使用支持事务的innodb引擎。但是,如果把数据库引擎变为myisam,那么程序也就不再支持事务了! ### Spring对事务的支持 (1)编程式事务管理 通过TransactionTemplate或者TransactionManager手动管理事务,实际应用中很少使用,但是对于你理解Spring事务管理原理有帮助。 (2)声明式事务管理 推荐使用(代码侵入性最小),实际是通过AOP实现(基于@Transactional的全注解方式使用最多)。 ![](https://pic.code-nav.cn/post_picture/1734931576698986498/YWmCHpJHLPmc1itU.webp) ### 声明式事务在多线程场景下的问题 ![](https://pic.code-nav.cn/post_picture/1734931576698986498/MHnvmlO7BZJMf020.webp) ![](https://pic.code-nav.cn/post_picture/1734931576698986498/cxDIuSFdlW2N7ZBS.webp) ![](https://pic.code-nav.cn/post_picture/1734931576698986498/gA7vFgkShodB92ii.webp) ### 事务的属性 事务属性包含了5个方面:隔离级别、传播行为、回滚规则、是否只读、事务超时 **只读模式** @Transactional(readOnly = false) 只读模式可以提升查询事务的效率,推荐事务中只有查询代码时,使用只读模式。默认是false,一般情况下,都是在类上添加@Transactional注解,针对类下的查询方法可以通过再次添加@Transactional注解,设置为只读模式,从而提高查询的效率(因为查询并不会改变数据库的数据,所以本身就不需要事务)。 **超时时间** @Transactional(timeout = 3) 默认值是-1,代表永远不超时,设置timeout = 时间(秒数),超过时间,就会回滚事务和报异常(TransactionTimedOutException),如果类上设置了,方法也设置了事务注解,方法上的注解会覆盖掉类上的注解! **指定事务异常回滚规则** @Transactional(rollbackFor=Exception.class, noRollbackFor=FileNotFoundException.class) 默认只针对运行时异常回滚(error也会回滚),编译时异常不回滚。可以指定: rollbackFor属性:指定哪些异常类才会回滚,默认是RuntimeException and Error 异常方可回滚; noRollbackFor属性:指定哪些异常不会回滚,默认没有指定,如果指定,应该在rollbackFor的范围内! 为了让发生所有异常都进行事务的回滚,我们可以指定Exception异常来控制所有异常都回滚!即:rollbackFor = Exception.class **事务隔离级别** _@Transactional(isolation = Isolation.READ_COMMITTED)_ Spring支持的隔离级别枚举: Isolation.DEFAULT:使用数据库默认隔离级别。 Isolation.READ_UNCOMMITTED:读未提交。 Isolation.READ_COMMITTED:读已提交。 Isolation.REPEATABLE_READ:可重复读。 Isolation.SERIALIZABLE:串行化。 数据库事务的隔离级别是指在多个事务并发执行时,数据库系统为了保证数据一致性所遵循的规定。常见的隔离级别包括: * 读未提交(Read Uncommitted):事务可以读取未被提交的数据,容易产生脏读、不可重复读和幻读等问题。实现简单但不太安全,一般不用。 * 读已提交(Read Committed):事务只能读取已经提交的数据,可以避免脏读问题,但可能引发不可重复读和幻读。(大多数数据库的默认隔离级别) * 可重复读(Repeatable Read):确保在同一事务中多次读取同一数据时,结果一致,不管其他事务对数据做了什么修改。可以避免脏读和不可重复读,但仍有幻读的问题。(MySQL的默认隔离级别) * 串行化(Serializable):最高的隔离级别,完全禁止了并发,只允许一个事务执行完毕之后才能执行另一个事务。可以避免以上所有问题,但效率较低,不适用于高并发场景。 **事务传播行为** _@Transactional(propagation = Propagation.REQUIRED)_ 事务传播行为定义了**多个事务方法相互调用时,事务应该如何传播**。Spring提供了7种事务传播行为,通过@Transactional注解的propagation属性进行配置。 | 传播行为类型 | 说明 | | --- | --- | | REQUIRED(默认) | 如果当前存在事务,则加入该事务;如果当前没有事务,则创建一个新事务,保证最后仅有一个事务。 REQUIRES_NEW | 无论当前是否存在事务,都创建一个新事务,并挂起当前事务(挂起事务的意思是在当前事务执行过程中,暂时将其暂停,并开启一个新的事务;挂起事务后,当前事务的状态会被保存,直到新事务执行完毕后再恢复),最后会有多个独立的事务,所以如果后面的事务报异常了,前面事务已经修改的数据并不会跟着一起回滚,适用于独立事务的场景,比如日志记录。 SUPPORTS | 如果当前存在事务,则加入该事务;如果当前没有事务,则以非事务方式执行,适用于查询方法,不需要强制事务。 NOT_SUPPORTED | 以非事务方式执行操作,如果当前存在事务,则挂起该事务,适用于不需要事务支持的操作。 MANDATORY | 如果当前存在事务,则加入该事务;如果当前没有事务,则抛出异常,强制要求必须开启事务。 NEVER | 以非事务方式执行操作,如果当前存在事务,则抛出异常,强制要求不能开始事务。 NESTED | 如果当前存在事务,则在嵌套事务内执行;如果当前没有事务,则创建一个新事务,适用于需要部分回滚的场景。 | ### 事务失效的场景 1. 数据库引擎不支持事务,比如mysql数据库使用的是MyISAM引擎; 2. 事务注解所在的方法是非public的,因为SpringAOP默认使用CGLIB代理,无法给非public的方法创建代理,导致事务切面无法切入; ![](https://pic.code-nav.cn/post_picture/1734931576698986498/sPA8bWcCQfHjKI2b.webp) 1. 异常类型不正确或被捕获,默认只回滚运行时异常和error,受检异常(IO异常,SQL异常)被视为业务异常,默认会提交事务;另外事务代理会通过捕获目标方法抛出的异常来决定是回滚还是提交,如果异常在方法内部被catch捕获且没有被重新抛出,代理会认为方法执行成功,从而提交事务; 2. 一个没有事务注解的方法调用了同一个类中有注解的方法,这种情况下使用的是this对象,是真实的对象,而不是代理对象; 3. 错误配置了事务传播属性,比如配置成以非事务方式运行,事务也会不生效; 4. 试图在final或static方法上使用事务,CGLIB代理是通过生成目标类的子类来实现的,它无法重写final和static方法。 05 Spring设计模式 ------------- ### 工厂模式 Spring使用工厂模式可以通过BeanFactory或ApplicationContext创建bean对象。 **两者对比:** BeanFactory:延迟注入(使用到某个bean的时候才会注入),相比于ApplicationContext来说会占用更少的内存,程序启动速度更快。 ApplicationContext:容器启动的时候,不管你用没用到,一次性创建所有bean。BeanFactory仅提供了最基本的依赖注入支持,ApplicationContext扩展了BeanFactory, 除了有BeanFactory的功能还有额外更多功能(比如:支持基于观察者模式的事件驱动编程,另外Java EE的标准注解,如@Resource也自动支持),所以一般开发人员使用ApplicationContext会更多。 ### 单例设计模式 在我们的系统中,有一些对象其实我们只需要一个,比如说:线程池、缓存、日志对象等。事实上,这一类对象只能有一个实例,如果制造出多个实例就可能会导致一些问题的产生,比如:资源使用过量、或者结果不一致性。 **使用单例模式的好处**: 对于频繁使用的对象,可以省略创建对象所花费的时间,这对于那些大对象而言,是非常可观的一笔系统开销; 由于new操作的次数减少,因而对系统内存的使用频率也会降低,这将减轻GC压力,缩短GC停顿时间。 **Spring中bean的默认作用域就是singleton(单例)的。** 除了singleton作用域,Spring中bean还有下面几种作用域: * **prototype**:每次获取都会创建一个新的bean实例。也就是说,连续getBean()两次,得到的是不同的Bean实例。 * **request**(仅Web应用可用):每一次HTTP请求都会产生一个新的bean(请求bean),该bean仅在当前HTTP request内有效。 * **session**(仅Web应用可用):每一次来自新session的HTTP请求都会产生一个新的bean(会话bean),该bean仅在当前HTTP session内有效。 * **global-session**(仅Web应用可用):每个Web应用在启动时创建一个Bean(应用Bean),该bean仅在当前应用启动时间内有效。 ### 代理模式 **一个经典的例子就是AOP,** 能够将那些与业务无关,却为业务模块所共同调用的逻辑(例如事务处理、日志管理、权限控制等)封装起来,便于减少系统的重复代码,降低模块间的耦合度。 Spring AOP就是基于动态代理的,如果要代理的对象,实现了某个接口,那么Spring AOP会使用**JDK**去创建代理对象,而对于没有实现接口的对象,Spring AOP会使用**Cglib**生成一个被代理对象的子类来作为代理。 ### 观察者模式 观察者模式表示的是一种对象与对象之间具有依赖关系,当一个对象发生改变的时候,依赖这个对象的所有对象也会做出反应。Spring事件驱动模型就是观察者模式很经典的一个应用。比如我们每次添加商品的时候都需要重新更新商品索引,这个时候就可以利用观察者模式来解决这个问题。 ### 适配器模式 适配器模式(Adapter Pattern)可以将一个接口转换成我们希望的另一个接口,适配器模式使接口不兼容的那些类可以一起工作。 # RabbitMQ部分 01 RabbitMQ的底层架构 ---------------- ![](https://pic.code-nav.cn/post_picture/1734931576698986498/pmWJVvFqIw7n1XCB.webp) 02 消息什么情况下会进入死信队列? ------------------ ![](https://pic.code-nav.cn/post_picture/1734931576698986498/2Pzh38RoCaYzjzoH.webp) 03 RabbitMQ怎样实现延迟队列? -------------------- ![](https://pic.code-nav.cn/post_picture/1734931576698986498/gh05HXfvlKne74x0.webp) 04 RabbitMQ中无法路由的消息会去哪里? ------------------------ ![](https://pic.code-nav.cn/post_picture/1734931576698986498/50mkHs7WwXVEbK1D.webp) 05 如何避免重复处理消息? -------------- ![](https://pic.code-nav.cn/post_picture/1734931576698986498/9rr0CzAmtyQxhV4M.webp) 06 RabbitMQ的推和拉模式? ------------------ 推模式: 推模式也称为订阅模式。消息是主动推送给消费者的,消费者会预先注册一个回调函数(消费者处理器),消息到达队列后会立刻推送给消费者,消费者设置预取数量来控制流量。这样的话消息的实时性高,到达后立即推送,减少了不必要的请求开销;缺点的话可能因处理能力不足导致消息堆积,需要合理设置预取数量以避免过载。适用于实时性要求高,消息量稳定的场景。 拉模式: 拉模式需要消费者主动从队列中获取消息。消费者主动请求,可以精确控制获取消息的时机和数量,适用于消息处理耗时较长或需要批量处理的场景。消费者可以按照自身能力获取消息,避免消息积压在消费者端,适合处理耗时任务。但是实时性较差,需要轮询去获取消息,大量的获取请求浪费可能增加了网络请求开销。适用于处理耗时,需要精确控制的场景。 07 消费者消费失败的常见因素? ---------------- (1)消息处理逻辑错误:比如消息内容是订单支付成功通知,但是消费者处理时因逻辑错误(比如金额计算错误)导致异常;也可能是外部依赖的异常,比如消息需要调用第三方的API,但是接口返回超时或者错误响应 (2)消费者也可能因为服务突然宕机,导致正在处理的消息未确认;或者消费者处理消息时执行耗时操作(比如生成大型图表),长时间未返回ACK触发RabbitMQ的超时机制。 (3)网络或中间件问题:消费者与RabbitMQ之间的网络抖动,导致心跳超时或通道(Channel)关闭;或者集群中某个节点宕机,消费者未正确切换到其他节点,导致消息无法投递。 08 消费者怎么保证消息消费的可靠性? ------------------- ![](https://pic.code-nav.cn/post_picture/1734931576698986498/z4yLyUf2tuL1uepC.webp) 09 RabbitMQ消息丟失的3种情况和对应解决? -------------------------- ![](https://pic.code-nav.cn/post_picture/1734931576698986498/9qzzoOy2mnhlk0tm.webp) (1)生产者弄丢数据 生产者将数据发送到RabbitMQ的时候,可能数据在半路弄丢了(比如网络问题之类的)。RabbitMQ有两种解决方式:1、生产者发送数据之前开启RabbitMQ的**事务功能**,如果消息没有成功被MQ接收到,那么生产者就会收到异常报错,此时就可以回滚事务,然后**重试发送消息**,但是问题是事务机制是同步的,提交一个事务之后就会阻塞等待事务处理完,所以吞吐量会下降,比较消耗性能; 2、开启MQ的**confirm模式**,这样每次写的消息都会分配一个唯一的id,然后如果写入了RabbitMQ中,MQ就会回传一个ack消息,告诉生产者这个消息ok了,如果MQ没能处理这个消息就会回调生产者的nack接口,告诉生产者消息接收失败了,然后可以进行重发。这个过程是异步的,发送完这个消息之后可以紧接着发生下一个消息。 (2)RabbitMQ弄丢了数据 这个我们需要开启RabbitMQ的持久化,保证消息写入之后会**持久化**到磁盘,这样即使RabbitMQ自己挂了,恢复之后也可以自动读取之前存储的数据。(需要在创建队列的时候给queue设置为持久化,并且还需要给消息也设置为持久化的) (3)消费者弄丢了数据 比如消息刚收到消费者就宕机了,这种情况我们需要**关闭MQ的默认ack机制**(因为默认是收到了就会触发ack确认),我们应该在消费者这边**业务处理完成之后再手动去ack确认**,这样即使消费者宕机没能处理消息,消息也不会被ack,所以消息可以重新入队处理。 10 RabbitMQ消息怎么传输? ------------------ ![](https://pic.code-nav.cn/post_picture/1734931576698986498/nRxNO1tQ5PXbUDLR.webp) 11 RabbitMQ怎么保证消息的顺序性? ---------------------- **方案一:单队列单消费者(FIFO最简单模式)** 这是最直接但也限制最大的方法。 **原理**: * **生产者**:将需要保证顺序的所有消息发送到**同一个队列**中。 * **消费者**:该队列**只能有一个消费者**,并且设置channel.basicQos(1),即每次只取一条消息,处理完一条再取下一条。 **优点**: * 实现简单,充分利用了RabbitMQ队列本身的FIFO(先进先出)特性。 **缺点**: * **无法水平扩展**:单个消费者是性能瓶颈,吞吐量低。 * **单点风险**:如果该消费者宕机,虽然消息不会丢,但处理会完全停止。 * **适用场景**:消息量非常小,对吞吐量要求不高的场景。 **方案二:根据消息ID或业务键进行分组(路由到同队列)** 这是最常用且合理的方案,核心思想是:**将需要保证顺序的消息通过路由键确保它们进入同一个队列,并被同一个消费者顺序处理**。 **原理**: * **消息分组**:对消息进行分组。例如,订单ID为order_123的所有消息(创建、付款、发货)必须顺序处理。 * **生产端**:使用**一致性哈希交换器**或**自定义路由键**,将同一组的消息总是路由到同一个队列。 **例如**:使用x-modulus-hash这类交换器,以订单ID order_123作为路由键,计算出的哈希值总会将其路由到队列queue_2。 * **消费端**:为**每个队列启动一个消费者**(可以是多个队列,即多个消费者实例,每个实例处理不同组的消息)。同样,每个消费者需要设置prefetchCount=1。 **工作流程**: 订单A(ID=1)的所有消息 → 路由键order_1→始终发往**队列1**→ 由**消费者1**处理。 * 订单B(ID=2)的所有消息 → 路由键order_2→ 始终发往**队列2**→ 由**消费者2**处理。 **优点**: * **高性能且可扩展**:不同的组(如不同的订单)可以被不同的消费者并行处理,解决了方案一的瓶颈问题。 * **逻辑清晰**:符合大部分业务场景(如订单、会话跟踪)。 **缺点**: * 需要提前规划好消息的分组逻辑。 * 如果某个组的消息特别多(“热点订单”),对应的队列和消费者可能会有压力,但通常这种情况较少。 12 RabbitMQ消息堆积怎么处理 ------------------- ![](https://pic.code-nav.cn/post_picture/1734931576698986498/24CrvdnhXWdkizT4.webp) ![](https://pic.code-nav.cn/post_picture/1734931576698986498/gmCvql1nPUgDqwvJ.webp) 13 优惠劵场景,RabbitMQ保证一致性问题 ------------------------ 本地消息表的核心思想是**将分布式事务拆分为两个本地事务**,通过数据库的事务特性和重试机制保证最终一致性: 第一个本地事务:在业务数据库中同时完成订单状态更新和消息记录 第二个本地事务:通过定时任务将记录的消息可靠地发送到MQ **具体实现流程:** **第一阶段:业务处理与消息记录** **开启数据库事务**:当优惠券被抢后,系统开始处理优惠券状态更新 **更新订单状态**:在同一个数据库事务中,将优惠券状态从"待发放"改为"已发放" **写入本地消息表**:在同一事务中,向专门设计的消息表插入一条记录,包含: 消息ID(唯一标识)、消息内容(JSON格式的优惠券数据)、消息状态(初始为"待发送")、创建时间、重试次数(初始为0) **提交事务**:只有当优惠券更新和消息记录都成功后才提交事务 **第二阶段:消息发送与状态更新** **定时任务扫描**:系统有一个独立的定时任务,定期(如每5秒)扫描本地消息表中状态为"待发送"的记录 **发送MQ消息**:对于每条待发送记录:尝试将消息内容发送到消息队列。如果发送成功:更新该消息记录状态为"已发送";如果发送失败:增加重试次数、记录错误日志、保持状态为"待发送" **重试机制**:对于发送失败的消息,定时任务会在下次扫描时重新尝试发送,直到:发送成功或达到最大重试次数(如5次),此时可将状态改为"发送失败"并报警。

消息队列 - 理论梳理

本文系统介绍消息队列概念,以及 RabbitMQ , RocketMQ , kafka 三个消息队列的核心概念。 ## 为什么需要消息队列? 在单体项目中,如果没有消息队列,那么: 1. 上游(用户)的操作(如生成视频)需要系统响应很久,那么上游就会长时间等待响应,不能做别的操作,阻塞线程。 2. 大量用户使用生成视频功能,每个任务下游直接处理(因为中间没有任何缓冲),可能导致下游服务处理不过来崩溃。 3. 系统可能在高峰时间段流量会突增,大部分时间处理的过来,如果盲目加机器,会造成成本效益低。 在分布式系统中,不同服务或模块之间需要通信。同步调用的缺点如下: 1. **系统耦合度过高**:服务间直接依赖,就像用胶水粘在一起。一个服务的变更或故障,可能直接影响到其他服务,维护和扩展变得困难。 2. **同步阻塞导致性能瓶颈**:主流程需要等待所有依赖操作完成才能继续。比如用户注册后,同步发送邮件、短信、初始化积分等,每一步都可能耗时,导致用户响应时间很长。 3. **无法应对突发流量(峰值冲击)** :在秒杀、大促等场景,瞬时流量远超系统处理能力。所有请求直接压到数据库等核心服务,极易导致系统崩溃。 **消息队列(MQ)** 的出现,就是为了解决这些痛点。它像是一个“中间人”或“缓冲带”,让服务间的通信更灵活、更可靠。 一般要用到消息队列多半是在分布式系统下。 ## 消息队列 ### 介绍 你可以把消息队列想象成一个**智能的邮局或消息中转站,临时(也可以持久)帮你存储待处理的任务。** - **消息**:便是待执行的任务,比如生成视频任务太耗时了,可以先放进消息队列。**消息本质上是一个自定义的 DTO对象**、JSON数据,并不是真的传递一个Task到消息队列,传递的是任务的描述信息,消费者接受后会根据信息执行任务。 - **生产者**:产生任务的上游,可以是用户,上游服务等等,生产者是一个抽象的概念,任何往消息队列放任务的角色都是生产者。 - **消费者**:处理消息队列内任务的服务,同样是抽象概念。 - **队列**:负责存储任务的容器,可以是 redis数据库,也可以用集合来手搓一个容器,最常用的是 rabbitMQ, RocketMQ , kafka ![img](https://pic.code-nav.cn/post_picture/1625164612255182850/0cPlvyWPKrZpguuj.webp) 消息队列有两个含义: **指代整个技术/系统:**当我们说:“我们项目里用了消息队列”,这里指的是一整套异步通信的技术和解决方案。 **指代具体的数据结构:**当我们说:“消息被发送到消息队列里”,这里特指那个存储消息的 FIFO 数据结构,也就是**队列**。 ### 消息队列的作用 消息队列的核心价值主要体现在五个方面 **解耦**:生产者和消费者无需直接知道对方的存在,只约定好消息格式即可。新增或减少一个消费者,生产者完全不用改代码,就像我们不用知道快递小哥的存在,我们只要知道菜鸟驿站在哪就行。 **异步**:生产者发送消息后立即返回,不用等待消费者处理。非核心操作(如发送邮件、记录日志)异步处理,大大提升主流程响应速度和用户体验 **削峰填谷**:在流量高峰期,消息队列充当缓冲区,暂存大量请求。后端服务按照自己的处理能力平稳地从队列中消费消息,避免被瞬间洪流冲垮 **可靠性:**消息队列一般会有持久化消息和防止消息丢失的功能。 **顺序性:**保证消息按顺序消费,主流消息队列都有顺序消费的实现。 ### 消息队列的模型 点对点(网上也说集群模式):一条消息只能有一个消费者消费,这条消息不能重复消费。 ![img](https://pic.code-nav.cn/post_picture/1625164612255182850/YaMFpTVNtA44GU1t.webp) 发布/订阅( 广播模型 ) : 消息队列会把一条消息广播给所有订阅了自己的消费者,都要老老实实处理。 ![img](https://pic.code-nav.cn/post_picture/1625164612255182850/V5vSqkRvR140ew6j.webp) ### 消息是什么? 我们知道生成者和消费者之间通过消息通信,那么消息是什么?是什么数据类型? 当我们说发送一个任务给消费者(服务端)执行的时候,发送的其实不是任务本身,而是一段任务描述或者执行任务所需的数据,服务端接受到任务的描述或者数据,就会在本地任务中执行。注意消息本身不是 Task 任务, 而是数据容器,某个数据类型。 无论你要发送什么数据类型,底层都会序列化成JSON 并转成 byte[] 数组传输,凡是能序列化JSON 或者 能转成 byte[] 的数据类型,都可以发送。 ```java Producer 任何数据类型 → JSON → byte[] → MQ → byte[] → Consumer JSON → 可选反序列化为Java对象 ``` 我们来看RocketMQ 能生产什么消息,看看如何发送。 String 类型消息: ```java public static void main(String[] args) throws Exception { // 声明一个默认的生成者 DefaultMQProducer producer = new DefaultMQProducer("example-producer-group"); // 绑定看板 producer.setNamesrvAddr("127.0.0.1:9876"); producer.start(); // 声明一条 String 类型消息 String strMsg = "Hello RocketMQ"; // 消息转成 byte[] 了 Message msg1 = new Message("TopicTest", "TagA", strMsg.getBytes("UTF-8")); producer.send(msg1); producer.shutdown(); } ``` 自定义对象转成 JSON ,作为一条信息: ```java class OrderDTO { String orderId; int count; } OrderDTO dto = new OrderDTO("A001", 10); byte[] body = JSON.toJSONString(dto).getBytes(StandardCharsets.UTF_8); // byte[] body = JSON.toJSONBytes(dto); 这样也行 producer.send(body); ``` Map / List /Set 等java集合: ```java Map<String, Object> map = new HashMap<>(); map.put("a", 1); map.put("b", "Hello"); byte[] body = JSON.toJSONBytes(map); producer.send(body); ``` 发送 byte[] 本身: ```java byte[] body = new byte[]{1,2,3,4}; producer.send(body); byte[] body = Files.readAllBytes(Paths.get("test.png")); producer.send(body); ``` 发送 XML : ```java String xml = "<user><id>1</id><name>Tom</name></user>"; byte[] body = xml.getBytes(StandardCharsets.UTF_8); producer.send(body); ``` 发送 CSV : ```java String csv = "id,name\n1,Tom"; byte[] body = csv.getBytes(StandardCharsets.UTF_8); producer.send(body); ``` 发送加密内容: ```java byte[] plain = "secret".getBytes(StandardCharsets.UTF_8); byte[] encrypted = encrypt(plain); producer.send(encrypted); ``` 发送压缩数据(gzip): ```java byte[] original = "Hello MQ".getBytes(StandardCharsets.UTF_8); ByteArrayOutputStream bos = new ByteArrayOutputStream(); GZIPOutputStream gzip = new GZIPOutputStream(bos); gzip.write(original); gzip.close(); byte[] body = bos.toByteArray(); producer.send(body); ``` 发送 Kryo 序列化后的对象: kryo 是一个序列化库,将java对象转换为可以网络传输的对象,不用Kryo也可以用Java原生序列化,别的序列化库。 ```java Kryo kryo = new Kryo(); // 创建 Kryo 实例 ByteArrayOutputStream bo = new ByteArrayOutputStream(); Output output = new Output(bo); // 绑定输出流 kryo.writeObject(output, dto); // dto 对象序列化成二进制,写入 output output.close(); // 关闭流 byte[] body = bo.toByteArray(); // 得到最终 byte[] mq.send(body); // 发送到消息队列 ``` **可以发送的消息包括:** - String - JSON - DTO / POJO - Map / List - 原始 byte[] - 文件(图片、视频、PDF、zip) - XML - CSV - Kryo 序列化对象 - Java 原生序列化对象 - 加密后的数据 - 压缩数据 本质上是参数的跨服务传递。 ### Rabbit MQ 消息队列概念汇总 RabbitMQ 和 RocketMQ 关键概念区别:RabbitMQ没有 **NameServer** **,** 使用exchange + binding key + routing key 实现路由。 RabbitMQ 的集群信息和 Broker 数量由集群内部自己维护,内部是Erlang写的,Erlang自带一个轻量级数据库,维护节点元信息,所以不需要 NameServer 。 **生产者**:生产消息给队列,生产者绑定交换机,只能往交换机生产消息,生产者发布消息只能指定一个交换机,路由键可以是通配符。 **交换机**(**exchange** ):交换机通过路由键和绑定键的匹配,将生产者的发来的消息路由到队列中。有多种类型交换机。 ![img](https://pic.code-nav.cn/post_picture/1625164612255182850/ULulRNVHkhjAWwJV.png) **消费者**:消费队列的消息,没有消费者组的概念,消费者只能绑定队列,一个消费者可以监听多个队列。 **队列(queue)**:存储消息的容器,Broker 管理的就是队列,通过绑定建,将队列绑定到交换机,可以声明绑定多个交换机。 **routing key** : 路由键,在消息形成时指定路由键,路由键可以是通配符,一条消息发给多个队列。 **binding key** : 绑定键, 队列在绑定到 Exchange 时所设置的匹配规则, Exchange 根据 routing key 与 binding key 是否匹配,决定是否将消息投递给该队列。可以是通配符 **通配符规则:让消息的传递变的很灵活。** | 通配符 | 含义 | 示例 | | ------ | ----------------------- | ------------------------------------------------------------ | | `*` | 匹配 **恰好一个单词** | 绑定键:`china.*.weather` 匹配:`china.news.weather` ✅ 不匹配:`china.weather` ❌ | | `#` | 匹配 **零个或多个单词** | 绑定键:`china.#` 匹配:`china.news` ✅、`china.news.weather` ✅、`china` ✅ | **Connection** : TCP 连接,由客户端(Producer 或 Consumer)建立 ,一个 TCP 连接是单个通道,消息拆成多个包,那只能一个包一个包发送,而且只能先把这个消息的包发完,再发送下一个消息的包, 如果客户端多个线程同时往这个 TCP 发送消息,数据包会互相混在一起。如果开通多个Connection连接, 会多次 TCP 握手,消耗资源,于是 Channel 就来拯救这种情况。 **Channel(通道)**:在一个连接里可以开多个轻量级通道,通道负责发送/接收消息。 这是RabbiMQ做的优化,一个逻辑概念,不是物理通道,给消息绑定一个 Channel_id 和channel会话实现,这样一个TCP连接 + 多个Channel 实现了多线程同时往 TCP 写消息。 **channel会话**:一个逻辑概念,本质上是一个数据结构实现,一个会话对象存储了 : - **Channel 状态**:open / close - **事务状态**(transaction) - **消息确认状态**(哪些消息已经 ack/nack) - **QoS(prefetch)**:当前可以未 ack 的消息数 - **绑定的队列、交换机信息** 把会话对象和Channel_id 绑定在消息中,便实现了一个 TCP + 多个Channel , 这样多条消息的数据包可以混合、交替发送, 最大化利用单个 TCP 连接资源 ,榨干 TCP 。 #### RabbitMQ 属性(队列类型) rabbitMQ 的队列可以设置属性,让队列具备某中特性,属性可以组合搭配: **1)持久化队列(durable=true) :** Broker 重启后队列仍存在, Broker 在内存中维护队列对象,同时写入 **磁盘文件。** - 原理:消息设置 delivery_mode=2 ,便也支持持久化, Broker 收到消息后 , 将消息写入磁盘 , 仅当消息安全写入磁盘后才返回 ack 给 Producer 。 - 当 **消费者消费并确认(ack)消息** 后,Broker 才会把消息从队列和磁盘上删除。 使用场景: 关键业务消息 (电商订单,物流消息推送) , 任务队列 (视频生成,图片处理等人物)都需要可靠性。 **2)独占队列(exclusive=true):** 一种 只能被创建它的 Connection 使用的队列。 - 原理:当客户端(Connection)声明一个 **独占队列,** Broker 会在内存中为这个队列分配一个逻辑对象,也就是connection对象的引用,以此来维护队列和连接的关联性。 - - 队列一旦创建,就和创建它的 Connection 绑定。 - 当这个 Connection 关闭或断开时,队列会被自动删除(如果同时设置了 `auto-delete`)。 - 只有创建该队列的 Connection 可以声明、消费消息。 - 其他 Connection 无法访问(无法订阅或发送消息到这个队列)。 使用场景: 临时队列场景,在线客服系统 (每次打开客户咨询是新的对话), 多人协作软件、游戏实时状态更新 **3)自动删除队列(autoDelete=true):** 当 最后一个消费者断开后,队列会被 Broker 自动删除。 - 原理: Broker 为队列维护一个 **消费者计数,**当消费者全部断开, Broker 检查 `autoDelete=true`, 队列被自动删除 使用场景:和独占队列一样适用于临时队列,不同的是队列销毁方式。 **4)临时队列:**不指定队列名, Broker 自动生成一个唯一队列名 , 通常格式类似 `amq.gen-<随机字符串>` - 原理: 队列对象存储在 Broker 内存中,绑定以下属性: - - `exclusive=true` → 队列绑定创建它的 Connection - `autoDelete=true` → Connection 关闭时自动删除 使用场景:这是 匿名临时独占队列 ,临时队列使用场景都相似,多一种技术选型。 #### 消息特性 决定可靠性、消费顺序和优先级 | 概念 | 定义 | 作用 / 使用场景 | | ------------------------------- | ------------------------------------------------------------ | ------------------------------------------- | | **持久化(Delivery Mode=2)** | 消息写入磁盘,保证 Broker 重启后不会丢失 | 关键业务消息,如订单、支付、任务队列 | | **非持久化(Delivery Mode=1)** | 消息只存在内存,Broker 重启会丢失 | 临时通知、日志、实时数据 | | **消息确认(Ack / Nack)** | 消费者处理消息后向 Broker 确认,Broker 才删除消息 | 确保消息被正确消费;支持 at-least-once 语义 | | **预取 / QoS(prefetch)** | 每个消费者一次可以接收的未 ack 消息数量 | 控制消费速率,防止消费者处理不过来 | | **TTL(Time-To-Live)** | 消息过期时间,到期后自动删除 | 临时消息、延迟消息 | | **优先级消息** | 消息带优先级,消费者先消费高优先级消息,**只能发送给优先队列** | 异步任务队列、紧急事件处理 | | **死信(Dead Letter)** | 无法正常消费的消息(拒绝、过期、队列满)被转入 DLX | 记录异常消息,做后续补偿或监控 | #### 交换机(Exchange) 决定消息路由策略 | 类型 | 定义 | 消息路由方式 | 使用场景 | | ------------------------------------ | ------------------------------------------------ | -------------------------------------- | -------------------------- | | **Direct Exchange** | 直连交换机 | 根据 `routing key` 精确匹配队列 | 单点消息投递,例如任务队列 | | **Fanout Exchange** | 扇出交换机 | 广播消息到绑定的所有队列 | 广播通知、消息广播系统 | | **Topic Exchange** | 主题交换机 | 支持通配符 `*`、`#` 匹配 `routing key` | 日志系统、主题订阅 | | **Headers Exchange** | 头交换机 | 根据消息头字段匹配队列 | 灵活路由、条件匹配 | | **默认交换机(Default / nameless)** | 每个队列都绑定到默认交换机,`routing key=队列名` | 自动直连 | 简单单队列发送 | #### 消费模式 决定消息接收和确认策略 | 模式 | 定义 | 特点 / 使用场景 | | ----------------------------- | -------------------------------- | ---------------------------- | | **Push 模式** | Broker 主动推送消息给消费者 | 实时性好,常用模式 | | **Pull 模式** | 消费者主动拉取消息 | 控制消费节奏,可结合批量拉取 | | **自动 ack(auto-ack=true)** | 消费者接收消息后自动确认 | 简单快速,但存在消息丢失风险 | | **手动 ack(manual ack)** | 消费者处理完再发送 ack | 消费可靠性高,支持异常重试 | | **事务模式(tx)** | 发送或确认消息可回滚 | 确保发送原子性,但性能低 | | **Confirm 模式** | 生产者确认消息是否被 Broker 收到 | 高性能可靠发送,替代事务 | 如果把 RabbitMQ 的特性比作一个工具箱,事务就像是那个被放在角落里、布满灰尘、很少被想起来的旧工具, RocketMQ 才是使用事务的消息对立,事务能力很强。 #### 高级特性 支撑业务扩展和可靠性保障 | 特性 | 定义 | 使用场景 | | ------------------------ | -------------------------------------------------- | ------------------------ | | **死信队列(DLQ)** | 消息因拒绝、过期或队列满无法消费,进入指定死信队列 | 异常消息处理、补偿机制 | | **TTL(消息/队列过期)** | 消息或队列超过指定时间自动删除 | 延迟消息、临时消息 | | **优先级队列** | 队列支持消息优先级,高优先级消息先消费 | 紧急任务处理、异步调度 | | **镜像队列(HA Queue)** | 队列在集群节点间复制多份,保证高可用 | 集群高可用、故障恢复 | | **延迟队列 / 延时消息** | 消息延迟指定时间后再投递 | 定时任务、定时提醒 | | **批量确认 / 批量发送** | 对消息进行批量 ack 或批量发布 | 提高吞吐量,减少网络开销 | #### 生产者工作原理 生产者的核心任务是:**将消息安全、高效地发送到指定的 Exchange**。它不关心消息的去向,只关心发送动作本身。 1. **建立连接与通道**: - - 生产者首先与 RabbitMQ Broker 建立一个 **TCP 连接**。 - 在这个连接之上,它会创建一个或多个**通道**。通道是轻量级的虚拟连接,所有的操作都在通道中进行,避免了为每个操作都建立 TCP 连接的开销。 1. **开启可靠性模式(关键)**: - - 为了确保消息不丢失,生产者需要开启可靠性机制。有两种选择: - - - **事务模式**:通过 `**channel.txSelect()**` 开启。发送一批消息后,调用 `**channel.txCommit()**` 提交或 `**channel.txRollback()**` 回滚。**此模式性能极差,已不推荐使用**。 - **发布者确认模式**:通过 `**channel.confirmSelect()**` 开启。这是**异步**的高性能模式。Broker 在成功接收消息后,会异步地向生产者发送一个确认(ACK)。 1. **发布消息**: - - 生产者调用 `**channel.basicPublish()**` 方法发送消息。 - 此方法需要指定核心参数:**Exchange 名称**、**Routing Key**(路由键)和**消息体**(包含消息内容和各种属性)。 1. **接收确认**: - - 在发布者确认模式下,生产者无需阻塞等待。它可以通过监听器异步接收 Broker 返回的 `**basic.ack**`。 - 如果收到 ACK,表示消息已成功到达 Broker。如果未收到(或收到 `**nack**`),生产者可以选择重发。 ![img](https://pic.code-nav.cn/post_picture/1625164612255182850/GbAxH7ACRl3zOovm.svg) #### Exchange 工作原理 Exchange 是 RabbitMQ 的**消息路由核心**。它接收来自生产者的消息,并根据**路由规则**将消息投递到一个或多个队列中。它本身不存储消息。 1. **接收消息**:Exchange 接收生产者发来的消息,并解析出其中的 **Routing Key**。 2. **匹配绑定**:Exchange 会查看与自己绑定的所有队列,以及每个队列的 **Binding Key**(绑定键)。 3. **执行路由**:根据自身的**类型**和**路由算法**,决定将消息投递到哪些队列: - - **Direct Exchange**:将 Routing Key 与 Binding Key 进行**精确匹配**。完全匹配则投递。 - **Fanout Exchange**:忽略 Routing Key,将消息广播到所有绑定的队列。 - **Topic Exchange**:将 Routing Key 与 Binding Key 进行**模式匹配**。支持 `*****`(匹配一个单词)和 `**#**`(匹配零个或多个单词)通配符。这是最灵活的类型。 - **Headers Exchange**:不依赖 Routing Key,而是根据消息的 **headers** 属性进行匹配。 1. **投递消息**:将消息的副本发送到所有匹配的队列中。如果没有任何队列匹配,消息的行为取决于 Exchange 的配置(可以丢弃或返回给生产者)。 ![img](https://pic.code-nav.cn/post_picture/1625164612255182850/wpBveDghFl04UUBZ.svg) #### 消费者工作原理 消费者的核心任务是:**从指定的队列中获取消息,并可靠地处理它**。 1. **建立连接与通道**:与生产者相同,消费者也需要建立连接(Connection)和通道(Channel)。 2. **声明队列与绑定**: - - 为了健壮性,消费者通常会尝试声明它要消费的队列(如果队列不存在则创建)。 - 如果队列还未绑定到 Exchange,消费者也会执行绑定操作。 1. **订阅消息**: - - 消费者通过 `**channel.basicConsume()**` 方法向 Broker **订阅**一个队列。 - 这相当于告诉 Broker:“请持续将这个队列里的消息推送给我。” 这就是**推模式**,也是最常用的模式。RocketMQ 是长轮询拉取模式。 - TCP层面的心跳检测, AMQP协议内置心跳(`heartbeat` 参数),连接断开会立即检测到。 而不是应用层用ping来做心跳检测。 - brocker维护`queue → consumer` 映射表 ,能看到消费者的活性 1. **接收并处理消息**: - - 当队列中有新消息时,Broker 会通过通道**主动推送**给消费者。 - 消费者的回调函数被触发,接收到消息并开始执行业务逻辑。 1. **发送确认**: - - 消息处理完成后,消费者必须向 Broker 发送一个**确认**。 - **手动 ACK**:通过 `**channel.basicAck()**` 显式确认。这是**推荐做法**,可以确保消息在处理失败时不会被丢失(可以重新入队或进入死信队列)。 - **自动 ACK**:消息一发送给消费者就自动确认。性能好,但容易在消费者处理失败时丢失消息。 ![img](https://pic.code-nav.cn/post_picture/1625164612255182850/KMo4jqc1KPBAwkMo.svg) #### 死信队列工作原理 死信队列是 RabbitMQ **容错和监控**机制的核心组成部分。 死信队列本质上是一个普通的队列,从逻辑上它被用来存储“死掉”的消息。这里的“死”不是指消息本身损坏,而是指**消息因为某些原因无法被正常消费**。 一条消息在以下三种情况下会成为“死信”: 1. **消息被拒绝**:消费者使用 `**basicNack**` 或 `**basicReject**`,并且设置了 `**requeue=false**`(不重新入队)。 2. **消息过期**:消息在队列中存活时间超过了设置的 TTL(下面会讲)。 3. **队列已满**:队列达到了设置的最大长度,无法再存入新消息。 死信队列本身就是一个普通队列,它的特殊之处在于它是通过**死信交换机** 配置而来的。 **死信交换机(DLX, Dead Letter Exchange)** 是一种特殊的交换机,用来接收那些在正常队列中无法被消费的“异常消息”。 当消息在一个队列中变成了“死信”,RabbitMQ 会自动把这条消息转发到对应的 **死信交换机**,再由 DLX 路由到一个新的队列(称为 **死信队列 Dead Letter Queue, DLQ**) 在RabbitMQ中落地实现很简单,其实就是普通交换机什么也不用改,交换机名字用dlx.xx开头比较合理,绑定的队列也用dlx.xx开头,特殊之处在于正常队列配置一个map参数,指定死信交换机,指定死信路由键就行。 死信交换机: ```bash // channel.exchangeDeclare 就声明了一个交换机,名字叫做 dlx.exchange channel.exchangeDeclare("dlx.exchange", "direct"); // 声明一个队列,名字叫做dlx.queue,意为死信队列,本质上是普通队列 channel.queueDeclare("dlx.queue", true, false, false, null); // 队列绑定死信交换机 channel.queueBind("dlx.queue", "dlx.exchange", "dlx.key"); ``` 正常交换机指定死信交换机,这一步才是真正给死信交换机赋予含义: ```java Map<String, Object> args = new HashMap<>(); args.put("x-dead-letter-exchange", "dlx.exchange"); // 指定死信交换机 args.put("x-dead-letter-routing-key", "dlx.key"); // 指定死信路由键(可选) channel.exchangeDeclare("normal.exchange", "direct"); // 最后一个参数不是null就代表有死信队列 channel.queueDeclare("normal.queue", true, false, false, args); channel.queueBind("normal.queue", "normal.exchange", "normal.key"); ``` 工作原理图: ![img](https://pic.code-nav.cn/post_picture/1625164612255182850/CP0mGUgwH6DsRnhF.svg) **死信队列使用场景**:降级处理、延时队列(别把死信队列看成故障的队列,就是个正常的队列) #### TTL(TIME-To-Live) **生存时间** TTL,即**生存时间**。它可以为消息或队列设置一个过期时间,超过这个时间后,消息就会“死亡”。T TL 是 RabbitMQ **消息生命周期管理**的基础功能。 **TTL 的两种设置方式** 1. **队列 TTL**: - - 在创建队列时,设置 `**x-message-ttl**` 参数(单位:毫秒)。 - **效果**:所有进入该队列的消息都会继承这个过期时间。 - **特点**:一旦设置,队列中所有消息的过期时间都一样,无法为单条消息定制。 1. **消息 TTL**: - - 在发送每条消息时,设置 `**expiration**` 属性(单位:毫秒)。 - **效果**:只有这条消息拥有独立的过期时间。 - **特点**:可以为每条消息设置不同的过期时间,更加灵活。 **TTL 的一个重要“坑”** **消息过期后,并不会立即从队列中删除!** RabbitMQ 只有在两种情况下才会处理过期消息: 1. 消息**即将被消费者消费**时(到达队列头部),没到达队列头部那就不会清理,哪怕过期了也是存在。 2. RabbitMQ 的一个**惰性扫描线程**定期检查队列。 这意味着,如果队列积压严重,一条已经过期的消息可能还会在队列中存活一段时间,直到它被扫描到或到达队首。如果配置了 DLX,过期消息在被处理时就会变成死信。 工作原理图: ![img](https://pic.code-nav.cn/post_picture/1625164612255182850/3wSBiVHwIkcEJ6eI.svg) #### 优先队列工作原理 一个支持消息优先级的队列。当队列中有消息积压时,高优先级的消息会**优先于**低优先级的消息被消费者获取。 优先级队列让 RabbitMQ 具备了**按重要性处理消息**的能力。 工作原理: 1. **启用优先级**:在创建队列时,必须设置 `**x-max-priority**` 参数(这就成为**优先队列**了),定义该队列支持的最大优先级(例如 5,表示优先级范围是 0-5)。 2. **设置消息优先级**:发送消息时,在消息属性中设置 `**priority**` 字段(值必须在队列支持的最大优先级范围内)。 3. **内部排序**:当消息进入队列时,RabbitMQ 并不是简单地追加到队尾。它会将高优先级的消息**插入到队列中较低优先级消息的前面**。 4. **消费顺序**:消费者总是从队列头部获取消息,因此高优先级的消息会先被消费。 优先级队列结构图: ![img](https://pic.code-nav.cn/post_picture/1625164612255182850/M4TNYs2q9tNzIndr.svg) #### 镜像队列工作原理 镜像队列:一个普通的队列,其内容可以被**实时复制**到一个或多个其他 Broker 节点上。形成“一主多从”的架构。 - **主节点**:负责处理所有对队列的读写操作(生产者发送、消费者消费)。 - **从节点**:作为热备,会从主节点同步所有消息和状态,但不对外提供服务。 镜像队列是 RabbitMQ **传统的高可用(HA)解决方案**,用于保证在主节点故障时,队列中的消息不丢失,服务不中断。 ⚠️**已过时**:在 RabbitMQ 3.8+ 版本后,官方推荐使用功能更强大、设计更现代的**仲裁队列**来替代镜像队列。 工作原理: 1. **配置镜像策略**:管理员需要在 RabbitMQ 管理界面或通过命令行设置一个**镜像策略**。这个策略定义了哪些队列需要被镜像,以及镜像到哪些节点上。 2. **主从同步**:当一个队列被策略匹配后,它就成为镜像队列。 - - 所有生产者发送的消息,都会先写入主节点。 - 主节点会将消息**同步**给所有从节点。 - 只有当**所有从节点**都确认收到后,主节点才会向生产者发送确认(ACK)。这保证了数据的强一致性。 1. **故障转移**: - - 如果主节点宕机,RabbitMQ 集群会自动从从节点中**选举一个新的主节点**。 - 原来的从节点升级为主节点,开始对外提供服务。 - 这个过程对生产者和消费者是**透明**的,它们会自动重连到新的主节点。 #### 仲裁队列工作原理 仲裁(Quorum)队列是 RabbitMQ **3.8 版本后推出的新一代高可用队列类型**。它基于 **Raft 共识算法**实现,旨在提供一个**数据安全、强一致、配置简单**的队列解决方案,是官方推荐的镜像队列替代品。 目标是实现**高可用**的**消息队列**架构。 “Quorum”一词意为“法定人数”,这暗示了它的核心机制:**需要大多数节点同意**,操作才能成功。 **工作原理:基于 Raft 共识算法** 仲裁队列的核心是 Raft 算法在消息队列场景下的实现。我们可以把它想象成一个**民主委员会**来管理队列。 一个仲裁队列通常部署在奇数个 Broker 节点上(最少 3 个),每个节点上的队列副本扮演一个角色: - **Leader(领导者)**:**唯一**的领导者,负责处理所有来自生产者和消费者的请求(读写操作)。 - **Follower(跟随者)**:普通的委员,不对外提供服务。它们的主要工作是**复制** Leader 的所有操作,并投票。 - **Learner(学习者)**:观察员角色,只同步日志,不参与投票。用于在不影响性能的情况下增加副本数,用于灾备或读取扩展。 - **消息写入流程(核心)** 这是理解仲裁队列的关键,它保证了数据的强一致性。 1. **所有写操作必须通过 Leader**。 2. **消息不是写入就成功**,而是需要**超过半数**的节点(包括 Leader 自己)都确认写入日志后,才算真正“提交”。 3. **只有已提交的消息才能被消费者消费**。这确保了即使 Leader 宕机,新选举出的 Leader 也一定拥有所有已提交的消息,**数据零丢失**。 ![img](https://pic.code-nav.cn/post_picture/1625164612255182850/nf0lSL02iwk08nAf.svg) - **故障转移流程** **心跳检测**:Leader 会定期向所有 Follower 发送心跳。 **触发选举**:如果 Follower 在一段时间内没有收到 Leader 的心跳,它就认为 Leader 宕机,于是发起选举,将自己转为 Candidate(候选人)。 **投票选举**:Candidate 向其他节点请求投票。获得**超过半数**选票的节点成为新的 Leader。 **数据恢复**:新 Leader 上拥有所有已提交的消息,可以立即对外提供服务,整个过程**自动完成**,对客户端透明。 #### 延迟队列工作原理 延迟队列(延迟消息):一个队列,其中的消息不会立即被消费者消费,而是在等待一段指定的时间后,才变成可消费状态。 延迟队列是一个非常常见的需求,但**RabbitMQ 本身不直接提供延迟队列功能**。我们需要通过插件或巧妙的组合来实现。 **实现方式一:TTL + DLQ(经典组合)** 这是最常用、最巧妙的实现方式,不依赖任何插件。 1. **架构设计**: - - 死信交换机:创建一个**业务交换机**和**业务队列**(不设置 TTL)。 - 延迟交换机:创建一个**延迟交换机**和**延迟队列**(设置 TTL)。 - 将延迟队列绑定到延迟交换机,并设置**死信交换机为业务交换机。** 1. **工作流程**: 1. 1. 生产者发送消息到**延迟交换机**,并设置 TTL。 2. 延迟交换机将消息路由到**延迟队列**。 3. 消息在延迟队列中等待,直到过期,期间不会有任何消费者拉取。 4. 消息过期后,成为死信,被发送到**死信交换机**(即业务交换机)。 5. 业务交换机根据路由规则,将消息路由到**业务队列**。 6. 消费者从业务队列中消费消息。 **实现方式二:延迟插件** RabbitMQ 官方提供了一个 `**rabbitmq_delayed_message_exchange**` 插件,提供了更原生、更优雅的实现。 1. **安装插件**:在所有 RabbitMQ 节点上安装并启用该插件。 2. **使用**:插件会提供一个新的交换机类型 `**x-delayed-message**`。 3. **工作流程**: 1. 1. 生产者发送消息到 `**x-delayed-message**` 类型的交换机。 2. 在消息头中添加 `**x-delay**` 属性,指定延迟时间(毫秒)。 3. 交换机收到消息后,**不会立即路由**,而是将其保存在一个内部的 Mnesia 表中。 4. 插件的后台定时器会定期检查这些消息。 5. 当消息的延迟时间到达后,交换机才会像普通交换机一样,将消息路由到目标队列。 ![img](https://pic.code-nav.cn/post_picture/1625164612255182850/7eKTabKaqxXXA9ss.svg) #### 批量发送/批量确认工作原理 **批量发送**:**生产者**不是一次发送一条消息,而是将多条消息打包,一次性通过一个网络请求发送给 Broker。 **批量确认**:**消费者**不是处理完一条消息就发送一次 ACK,而是处理完一批消息后,发送一个 ACK,确认这批消息都已成功处理。 这是 RabbitMQ **性能优化**的关键手段,通过减少网络往返次数来提升吞吐量。 **批量发送原理** - **客户端实现**:批量发送主要是**客户端**的行为。RabbitMQ 的 Java 客户端(如 Spring AMQP)提供了 `**RabbitTemplate.convertAndSend**` 的批量版本。 - **网络效率**:将 N 条消息的 N 次网络请求,合并为 1 次网络请求,极大地减少了网络 RTT(往返时间)的开销。 - **权衡**:批量发送会增加客户端的内存使用,并且会带来一定的延迟(需要等一批消息凑齐或超时),空间换时间。 **批量确认原理** 1. **开启手动确认**:消费者必须关闭自动 ACK (`**autoAck=false**`)。 2. **处理消息**:消费者从队列中拉取一批消息(或 Broker 推送一批),并依次处理。 3. **发送批量 ACK**:当一批消息都处理成功后,消费者调用 `**channel.basicAck(deliveryTag, multiple=true)**`。 - - `**deliveryTag**`:这批消息中**最后一条**消息的标签。 - `**multiple=true**`:这是关键!它告诉 Broker:“请确认并删除**小于等于**这个 `**deliveryTag**` 的所有消息”。 1. **Broker 处理**:Broker 收到这个批量 ACK 后,会一次性将队列中这批已确认的消息全部清除。 ![img](https://pic.code-nav.cn/post_picture/1625164612255182850/FgoAWFYWmEonFOkn.svg) ### RocketMQ 消息队列概念汇总 RocketMQ 非常适合应用在微服务架构中,经常作为微服务的消息中间件,所以下面的概念以微服务视角去思考、代入才会更好理解。 我们来看下面这张图,它囊括了⁢ RocketMQ 相关的组件和角色‏。 ![img](https://pic.code-nav.cn/post_picture/1625164612255182850/bzd54XQyzVPQn4E0.webp) 核心概念一览表 | 概念 | 定义 | 对应 RabbitMQ 类比 | 说明 / 使用场景 | | -------------------------------- | ------------------------------------- | ------------------------------------------------- | --------------------------------------------- | | **生成者和消费者** | | | | | **Producer** | 消息生产者 | RabbitMQ Producer | 发送消息到 Topic / Queue | | **Producer Group** | 消息生产者群组 | | 事务,管理多个生产者 | | **Consumer** | 消息消费者 | RabbitMQ Consumer | 拉取或订阅消息 | | **Consumer Group** | 消费者组 | RabbitMQ Queue + 多消费者 | 一组消费者共享 Topic 消息,实现负载均衡或广播 | | **中间队列层概念** | | | | | **Broker** | 消息服务器 | RabbitMQ Broker | 消息存储和分发节点 | | **NameServer** | 元数据服务 | RabbitMQ 没有对应(Exchange/Binding+Cluster管理) | 集群路由和 Broker 元数据注册中心 | | **Topic** | 消息主题 | RabbitMQ Exchange(逻辑路由) | 消息的逻辑分类,Producer 发送消息到 Topic | | **Message Queue (队列)** | 物理队列 | RabbitMQ Queue | 存储消息的物理队列,由 Broker 管理 | | **commitLog** | RocketMQ 核心物理文件 | / | 存储所有消息的实际内容 | | **ConsumeQueue** | 队列的索引文件 | / | 一条记录映射一条commitLog 记录 | | **Message Queue Offset** | 队列偏移量 | RabbitMQ 没有明确概念 | 消费进度记录在 Broker 或 Consumer | | **消息层概念** | | | | | **Message** | 消息对象 | RabbitMQ Message | 包含 Body、Topic、Tag、Keys 等 | | **Tag** | 消息标签 | RabbitMQ RoutingKey | 消费端可做消息过滤 | | **Message Key** | 消息唯一标识 | RabbitMQ MessageId / header | 用于消息追踪和查找 | | **顺序消息** | 消息保证严格顺序消费 | RabbitMQ 默认 FIFO(Queue 内) | 消息队列分区内顺序处理 | | **RocketMQ 事务消息** | Producer 发送事务消息,最终提交或回滚 | RabbitMQ Transaction 或 Confirm 模式 | 保证分布式事务可靠性 | | **延迟消息 / Scheduled Message** | 消息延迟投递 | RabbitMQ TTL + 延迟插件 | 定时任务、延迟任务 | | **死信消息(DLQ)** | 消费失败或超过重试次数的消息 | RabbitMQ DLX | 异常消息处理 | 概念很多,最核心的是 commitLog 、Consume Queue 、 message key 、事务消息 #### 生产者和消费者 **1) Producer** : **消息的发送方**,负责把消息发送到 **Topic(逻辑分类),** 消息通过 Broker 进行存储,最终写入 CommitLog, Producer 可以选择不同的发送模式: - **同步发送(sync)**:等待 Broker ACK,保证消息可靠性 - **异步发送(async)**:回调方式,适合高吞吐量场景 - **单向发送(one-way)**:不等待 ACK,适合日志或监控消息 发送消息可以对消息指定: - **指定 Topic 和 Tag** ,用于消息分类和过滤 。 - **指定 Message Key** ,业务唯一标识,便于查询或幂等处理 - **事务消息支持** , 可以发送事务消息,配合 Broker 做事务半消息/回查机制 **2) Producer Group(生产者群组)** : 每一个生产者实例在创建时,一定要指定一个**生产者组**名,因此形成生产者组的概念。作用是 事务消息回查的关键标识 —— Producer 发送半事务消息,Brocker总是等不到提交指令,就会根据生成者组名回查生产者现在是什么状态。 **3) 消费者**: 消费者是一个独立运行的程序或进程(比如一个 Spring Boot 微服务实例)。它的任务是从 RocketMQ 的服务器(Broker)获取消息,并执行具体的业务逻辑(比如扣减库存、发送短信、生成报表等)。 **4) Consumer Group(消费者组):**在消费者创建时必须指定一个消费者组名字,消费者组也就因此形成。作用是**处理同一类业务逻辑**的消费者应用实例会被归为同一个组,增大这一类业务的消费吞吐量。 - 集群模式(点对点):消费组会**分摊消费**Brocker 中的 Topic 消息,一个消费者实例对接一个Topic下的队列。 - 广播模式: 消费者组内的每一个消费者,都会收到该主题下的**全量**消息 。 consumer - Brocker 传递数据的过程: 1. 消费者主动向 mq 发送一个长轮询请求 2. 如果有数据立即返回 3. 如果没有数据,挂起请求15秒(消费端指定**长轮询**的时间,挂起过程不占用线程),期间内有数据,唤醒请求返回数据。期间内没有数据,返回空响应(空转)。 ------ #### 中间层队列概念 **1)Broker:** 消息存储与转发器,负责接收生产者的消息,接收的消息持久化存储到commitLog,Broker提供给消费者消息拉取能力。生成者通过NameServer 找到 Brocker **2)NameServer** : 元数据注册中心,管理 **Broker 地址** 和 **Topic 路由**信息 , 可以由轻量级服务器作为技术实现, 每个 Broker 启动时注册到 NameServer ,汇报broker的topic信息到看板并维持心跳检测。NameServer 会注册生产者组和消费者组的信息, **Producer 和 Consumer** 找NameServer 拉取 **Topic 路由信息并缓存。** **3)Topic (主题)**: 一个逻辑概念,由若干队列和Broker 绑定 topic 名字实现。 Topic 内部会被拆分成若干个 **队列**, 每个队列独立存储消息,形成并行的存储和消费结构。 Producer 和 Consumer 只需约定 Topic, 不关心彼此 ,即可互相通信, 实现解耦**。**Producer 通过 NameServer 查询 Topic 的分区信息 。 **4)message queue** : 存储消息的队列容器,消息在队列内有顺序编号(MessageQueue Offset),用于顺序消费 ,负责对topic **5)commitLog** : 是 RocketMQ 的 **核心物理存储文件**, 存储的就是 **消息的物理内容,** 它按顺序 顺序追加(append) 所有消息,无论消息属于哪个 Topic 或 MessageQueue, 所以整个RocketMQ 仅有一个 commitLog . commitlog = 消息内容 + 消息属性 + 系统元数据 + 长度信息 **6)Consume Queue** : 队列的元信息 , ConsumeQueue 是索引文件 ,每条记录对应 MessageQueue 的一个 Offset,记录了该逻辑消息在 CommitLog 中的物理偏移,Consumer 通过它快速找到消息 ,内部结构 : ```bash ConsumeQueue[0] = <CommitLog offset, size=128, tag> ConsumeQueue[1] = <CommitLog offset, size=128, tag> ConsumeQueue[2] = <CommitLog offset, size=128, tag> 这里的[0][1][2] 便是 Message Queue Offset ``` **7)Message Queue Offset :** 每一个消息中都有一个顺序编号,称为Offset , 这是 **物理队列内消息的唯一序号**,编号从0单调递增,假如队列内3条消息消费完了,再来一条新消息,从4开始编号, 不会重用 ,永远递增。 - RocketMQ 崩溃后如何恢复编号? 通过 commitLog + ConsumeQueue - 重启后, Broker 会 **读取 ConsumeQueue 文件尾部**, 找到最后一条记录 `<CommitLog offset, 消息长度, Tag>`, 逻辑序号 = 该记录在文件中的顺序位置 , 即使消息已经被 Consumer 消费,ConsumeQueue 索引仍存在,所以可以恢复最大 Offset . - **集群模式下,由 Broker 集中管理****Message Queue Offset,**确保一致性和可靠性。 - **广播模式下,由消费者各自本地管理****Message Queue Offset,**简单高效。 Rocket MQ 的 commitLog 是关键设计, commitLog 会保存已经消费的消息,不删任何消息,这样便支持重复消费,支持了重复消费,那么 同一个 Topic 就可以被多个 Consumer Group 消费 ,并以 commitLog 的设计为基础设计了事务消息。 注意: RocketMQ 并不是无限保留 CommitLog 消息 ,**消费完成并过期的消息**,会由 **定期清理(Clean-up)机制** 删除 ------ #### 消息层概念 **1) Message(消息):** 是 RocketMQ 中 **最基本的数据单元,** 表示 业务系统需要传递的事件或数据, 主要包含三个部分: - 主题(Topic) : 消息逻辑分类,消费者按 Topic 消费 - 消息体(Body) : 业务数据,通常是字节数组(例如 JSON、字符串) - 消息属性(Properties) : 元数据,可选,用于 Tag、Key、延迟、事务、过滤等 详细字段: | **字段** | **说明** | **用途** | | -------------- | ------------ | -------------------------------- | | Topic | 消息主题 | 消费者订阅、分类 | | Body | 消息内容 | 实际业务数据 | | Tag | 消息标签 | 消费端可选择性订阅,用于消息过滤 | | messge Key | 消息唯一标识 | 幂等、查询、追踪 | | DelayTimeLevel | 延迟级别 | 延迟消息发送 | | BornTimestamp | 消息产生时间 | 消息顺序/延迟计算 | | Properties | 扩展属性 | 自定义字段,可用于业务逻辑 | **2) Tag (标签)**: 理解成消息的一个字段,消息的元信息,动态的随消息生成。多个topic可以有相同的tag,不影响。 **3) message Key**: 是消息的唯一标识符或 业务唯一标识**,** 它是由生产者业务系统设置的标识 , 通常用于业务系统追踪或查询消息 , 可用于 **消息查询、幂等、日志追踪、统计,** Message Key 尽量唯一,但也可以多个消息共用一个 Key(例如同一个订单号) RocketMQ 存储消息时,Key 仅作为 **消息属性**,实际消息定位仍靠 **CommitLog + ConsumeQueue Offset** Message Key 可以唯一,也可以不唯一 - **唯一 Key** - - 适合 **业务系统要求严格幂等或精确查询**的场景 - 例如:订单号、支付流水号 - 查询时通过 Key 能唯一定位消息 - **非唯一 Key** - - 同一个业务对象可能发送多条消息,但 Key 可以相同 - 例如:同一个用户多次操作,Key 是用户 ID - 查询时可以返回多条消息 **4) 顺序消息:** 消息的发送顺序与消费顺序保持一致。 RocketMQ实现顺序消息的关键在于 **“同一业务的消息始终进入同一个 MessageQueue”**,然后在消费端 **同一时刻只由一个线程顺序消费该队列的消息**。 整个过程依赖队列内天然的FIFO. ------ ##### 1. 事务消息 **事务消息**:它确保 消息的发送 与 本地事务的执行 要么都成功,要么都失败,从而保证分布式场景下的数据一致性 本地事务指的是生产者服务开启了一个事务,将业务操作和发生消息都纳入到事务步骤中,同成功同失败。 ![img](https://pic.code-nav.cn/post_picture/1625164612255182850/Qp0sBuhBLlxpZFOe.svg) **半消息状态**:半消息会被标记一个特殊属性 `PROPERTY_TRANSACTION_PREPARED`,设置为 `"true"`,事务消息的独有特征,半消息不会直接发送到指定的业务topic中,RocketMQ 采用了 topic 隔离策略来存储半消息,这个 topic 也是个普通的topic, 但是消费者不能消费,便具有了半消息的语义。 **存半消息的 topic** :RocketMQ 内部名为 `RMQ_SYS_TRANS_HALF_TOPIC`的 Topic ,此 Topic 是内置的、天生的,作用是业务消费者无法订阅它。 **操作消息(OP消息)**:这是一个 RMQ_SYS_TRANS_OP_HALF_TOPIC,半消息提交或回滚的操作也叫做一条信息,存储在这里。因为 RocketMQ 的存储机制是基于 commitlog 顺序写 + ConsumeQueue 索引文件。 **消息一旦写入,就不会修改(Append-only 模式)**。 Broker 不能直接更新 Half Message 的状态, 所以采用“写一条对应操作记录”的方式来表达状态变化 ,**这和数据库的“WAL 日志(Write-Ahead Log)”很像**: 不改原数据,而是追加一条变更操作日志。 **事务回查**:rocketMQ 向 producer 询问当前的事务状态,producer 会检查本地事务是成功还是失败 **定时任务**: Broker 定期扫描 Half Queue,找出需要回查的消息 **Producer 回查接口** : 生产者实现业务逻辑检查点 在分布式系统中,最难的问题之一就是 **跨系统的一致性**,既 消息系统(MQ)和数据库事务无法保证**原子性(Atomicity)**。 RocketMQ 的 **事务消息** 正是为了解决这一点而生的。 它提供了一个“**两阶段提交 + 回查机制**”的模型,保证两者的最终一致性。 缺点:它只保证**最终一致性**,且中间状态(如支付成功但下游服务异常)在短暂时间内可能不一致,需要业务方自行处理。 Kafka、RabbitMQ 等常见 MQ 原生并没有类似机制 ,而这正是 RocketMQ 的亮点之一,国产之光! ------ ##### 2. 延迟消息 **延迟消息(Scheduled Message)** 是指消息在发送到 Broker 后,不会立即被投递给消费者,而是**延迟一段时间后**再被消费。 实现原理: - 生产者形成延迟消息时 , 带上 delayLevel,Broker 收到后 **不会立即写入真实 Topic** - Broker 写入“延迟队列”(定时队列) , 延迟队列的 Topic 固定为 `SCHEDULE_TOPIC_XXXX` - Broker 定时任务扫描到期消息 , 到期后将消息重新写入原始 Topic(即用户真正的 Topic) - 消费者收到消息 , 从原 Topic 消费到这条延迟后转发的消息 这跟RabbitMQ 的 **TTL + DLQ 实现延迟队列思想一致。** ------ ##### 3. 死信消息 **死信消息**:指在消息消费失败,且达到**最大重试次数**后(默认16次)仍未成功,RocketMQ 会将这类消息转移到特殊的队列中进行隔离,这些消息被称为死信消息,存储死信消息的队列就是死信队列 **死信队列**:它是一个**特殊的 Topic**,名称通常为 `**%DLQ% + 消费组ID**`。每个消费组都有其对应的死信队列,用于存储该消费组内所有 Topic 的死信消息 死信队列是 RocketMQ **消息可靠性设计的最后一道防线**,旨在**隔离问题消息**,防止其影响正常消息的处理流程,并提供一个集中的地方供监控和人工干预。 使用场景:死信队列主要用于**处理消费失败且无法通过重试恢复的消息** - **订单处理失败**:订单创建、支付、库存扣减等操作失败的消息,可以进入死信队列,后续通过人工干预或自动程序进行补偿 1. **异常消息处理**:由于数据格式错误、业务逻辑异常等原因无法消费的消息,进入死信队列,供后续排查和修复。 2. **超时消息处理**:设置了 TTL 但未能在规定时间内消费的消息。 3. **监控与告警**:监控死信队列可以及时发现系统异常,是系统健康度的重要指标。 ------ #### 生产者工作原理 单生产者工作原理 ![img](https://pic.code-nav.cn/post_picture/1625164612255182850/EV3VY9nAQH18AorY.svg) 1. **创建并启动 Producer** - - 程序创建 `**DefaultMQProducer**` 实例。 - 设置 NameServer 地址、Producer Group 名称。 - 调用 `**start()**` 方法启动 Producer。 1. **获取路由信息** - - Producer 启动后,会从**一个 NameServer**(随机选择)拉取 Topic 的路由信息。 - 路由信息包含:该 Topic 分布在哪些 Broker 上,每个 Broker 上有哪些 MessageQueue(队列)。 - Producer 会**定时(默认30秒)** 更新本地路由缓存。 1. **选择 MessageQueue** - - Producer 根据负载均衡策略(如轮询 `**RoundRobin**`),从获取到的 MessageQueue 列表中选择一个队列。 - 如果选择了写失败的队列,会进行重试,并可能暂时规避该队列。 1. **构建并发送消息** - - 将消息体、Topic、Tag、Key 等信息封装成消息对象。 - 对消息进行序列化、压缩等处理。 - 通过 Netty 客户端向选定的 Broker 发送消息。 1. **Broker 处理消息** - - Broker 接收消息后,将其**顺序写入 CommitLog 文件**。 - 写入成功后,返回一个包含消息物理偏移量(CommitLog Offset)的 ACK 给 Producer。 1. **Broker 异步构建索引** - - Broker 的后台线程会**异步**地将 CommitLog 中的消息位置、大小等信息分发到对应的 ConsumeQueue 文件。 - 如果消息带有 Key,还会构建 IndexFile 索引,以支持按 Key 查询。 1. **Producer 接收 ACK** - - Producer 收到 Broker 返回的 `**SendResult**`,状态为 `**SEND_OK**`,表示消息发送成功。 ![img](https://pic.code-nav.cn/post_picture/1625164612255182850/OQ81xszfP89OhFUU.svg) 重点说明: - **异步构建索引**:消息写入 CommitLog 和构建 ConsumeQueue 索引是**异步**的。这极大地提升了 Broker 的吞吐性能,因为写入 CommitLog 是顺序写,非常快,而构建索引可以稍后批量处理。 - **高可用设计**:如果发送失败,Producer 会自动重试(默认2次),并可能选择其他 Broker 上的队列,实现了容错。 **生产者群组工作原理** ![img](https://pic.code-nav.cn/post_picture/1625164612255182850/kY2SythhNRuqC085.svg) **1. 启动与路由发现** - **拉取路由**:当 Producer Group 中的任何一个实例启动时,它会**定期向 NameServer 拉取**它所关心的 Topic 的路由信息(包括该 Topic 分布在哪些 Broker 上,以及每个 Broker 上的 MessageQueue 列表)。 - **心跳机制**:Producer 实例会与 NameServer 和 Broker 保持心跳,但这主要是为了让它们感知自己的存活状态,而非为了“注册”。 - **无状态设计**:因此,Producer 实例本身是**无状态的**,Broker 也只关心消息内容,不关心消息具体来自哪个实例。这使得 Producer Group 内的实例可以任意增减,实现弹性伸缩。 **2. 消息发送(正常运行)** 在发送消息时,Producer Group 的工作机制体现了其高扩展性: - **独立发送**:每个 Producer 实例都**独立工作**,根据其本地缓存的路由信息,通过负载均衡策略(如轮询)从 Topic 的多个 MessageQueue 中选择一个。 - **并行处理**:然后,它将消息直接发送给该 MessageQueue 所在的 Broker。多个实例可以并行地向不同的 Broker 或不同的 Queue 发送消息,共同承担发送压力,实现了负载均衡。 **3. 事务消息回查(容错机制)** 在事务消息场景下,Producer Group 的设计提供了关键的容错能力,这也是其核心价值之一: - **问题场景**:如果某个 Producer 实例在发送“半消息”后、提交最终状态(COMMIT/ROLLBACK)前突然宕机,Broker 上就会留下一个状态不明的“孤儿”半消息。 - **回查机制**:此时,Broker **不会主动“通知”** 同组的其他实例。正确的流程是,Broker 会**主动发起回查**。它会根据半消息中记录的 **Producer Group 名称**,向 NameServer 查询该 Group 下**任意一个存活的 Producer 实例地址**。 - **状态恢复**:然后,Broker 向这个被选中的实例发起回查请求,由它代表整个 Group 去查询本地事务状态(如查询数据库),并将结果(COMMIT 或 ROLLBACK)返回给 Broker。Broker 根据这个结果完成最终的事务提交或回滚。 ![img](https://pic.code-nav.cn/post_picture/1625164612255182850/Lwx9C4KAJLy6MSX2.svg) ------ #### 消费者工作原理 在深入两种模式之前,先理解几个贯穿始终的核心概念: - **Consumer Group(消费者组)**:一个逻辑概念,由多个消费者实例组成。它们共同消费一个或多个 Topic,是实现负载均衡和容错的基本单元。理解消费者实例和队列对应关系很重要。 - **Rebalance(重平衡)**:一个动态过程,当消费者组内的实例发生变化(加入、离开)或 Topic 的队列数量变化时,Broker 会重新分配队列与消费者的对应关系。 - **Offset 管理**:记录消费进度。这是区分两种模式的关键所在。 **1)集群模式** 集群模式是**默认且最常用**的模式,旨在通过横向扩展来提升消费能力。 1. **队列分配**:一个 Topic 的多个 MessageQueue 会被**分配**给消费者组内的不同实例。**一个队列在同一时间只会被组内的一个消费者实例消费**。 2. **负载均衡**:通过 Rebalance 机制,如果消费者实例数量多于队列数量,多余的实例会闲置;如果实例数量少于队列数量,一个实例会消费多个队列。这实现了消费任务的负载均衡。 3. **故障转移**:如果一个消费者实例宕机,它负责的队列会被 Rebalance 机制重新分配给其他存活的实例,确保消息消费不中断。 4. **Offset 存储**:由于消费进度需要被组内所有实例共享(尤其是在故障转移时),消费进度(Offset)**集中存储在 Broker 端**。 ![img](https://pic.code-nav.cn/post_picture/1625164612255182850/DxRIkMpYBRki7p6P.svg) **2)广播模式** 广播模式适用于需要将同一条消息推送给所有下游消费者的场景,如配置更新、状态通知等。 1. **全局消费**:消费者组内的**每个实例都会消费 Topic 下所有 MessageQueue 的所有消息**。消息不会被分摊,而是被广播。 2. **无 Rebalance 分配**:由于每个消费者都要消费所有队列,因此不存在队列分配的 Rebalance 过程。每个消费者独立地订阅所有队列。 3. **独立消费**:各消费者实例的消费进度互不影响,一个实例消费慢或失败,不影响其他实例。 4. **Offset 存储**:由于每个消费者的消费进度都是独立的,不需要与其他实例同步,因此消费进度(Offset)**存储在消费者本地**。 ![img](https://pic.code-nav.cn/post_picture/1625164612255182850/3XjuBnNUWeOBseDP.svg) 两种模式对比: | 特性维度 | 集群模式 | 广播模式 | | --------------- | ---------------------------------------------------- | -------------------------------------------------------- | | **消费关系** | 一个队列**只被**一个消费者实例消费 | 一个队列**被所有**消费者实例消费 | | **负载均衡** | **支持**,通过 Rebalance 在组内分摊消费任务 | **不支持**,每个实例承担全部消费负载 | | **故障转移** | **支持**,宕机实例的队列会被其他实例接管 | **不支持**,实例间完全独立,互不影响 | | **Offset 存储** | **Broker 端**集中存储 | **消费者本地**独立存储 | | **适用场景** | 高吞吐量、分摊消费压力的业务(如订单处理、日志分析) | 配置下发、状态通知、所有客户端需要同步收到相同消息的场景 | ------ ### kafka 消息队列概念汇总 | 类别 | 概念 | 作用 | 类比与说明 | | -------------------- | -------------------------------- | ------------------------------------------------------------ | ------------------------------------------------------------ | | **核心组件** | **Producer** | 消息生产者,负责发送消息到 Kafka Broker。 | 消息的来源,可以是网站前端、后端服务、日志采集器等。 | | | **Consumer** | 消费者,从 Kafka 读取消息。 | 消息的终点,负责处理消息,如数据分析、业务处理等。 | | | **Broker** | Kafka 服务器实例,负责存储和转发消息。 | 一个 Kafka 集群由多个 Broker 组成,每个 Broker 都是一个独立的节点。 | | | **Topic** | 消息分类主题,是逻辑上的消息分类。 | 类似数据库中的表名,用于区分不同类型的消息(如订单、用户行为)。 | | | **Partition(queue)** | Topic 下的分区,是真正存储消息的地方。partition 就是**队列** | 一个 Topic 可以分为多个 Partition,分布在不同 Broker 上,实现水平扩展和并行处理。**类比 RocketMQ 的 MessageQueue**。 | | | **Offset** | 每条消息在 Partition 中的唯一序号。 | **Partition 内部**消息的指针,从 0 开始单调递增。**类比 RocketMQ 的 Message Queue Offset**。 | | | **Consumer Group** | 一组消费者,负责负载均衡消费。 | 一个 Topic 的消息可以被多个不同的 Consumer Group 订阅,每个 Group 都有独立的消费进度。 | | | **Leader / Follower** | Partition 的主从副本,保证高可用。 | 每个 Partition 有一个 Leader 负责读写,多个 Follower 负责同步数据。**类比 RocketMQ 的 Master/Slave**。 | | | **Zookeeper / Kafka Controller** | 管理 Broker、Topic、Partition 元数据。 | **旧版依赖 Zookeeper**。**新版(KRaft)用 Controller 代替**,Controller 是从 Broker 中选举出来的,简化了架构。**类比 RocketMQ 的 NameServer**。 | | **生产者端** | **acks** | 生产者可靠性核心配置,决定何时认为消息发送成功。 | `**acks=0**` (发完即成功), `**acks=1**` (Leader收到即成功), `**acks=all**` (所有副本收到才成功)。 | | | **幂等性生产者** | 保证单分区单会话内 Exactly-Once,防止重试导致消息重复。 | 开启后,Producer 为每条消息分配唯一ID,Broker 去重,确保消息不重复。 | | | **事务性生产者** | 保证跨分区、跨主题的原子写入,实现“要么都成功,要么都失败”。 | 结合幂等性和事务协调器,实现 Kafka 的 Exactly-Once 语义 (EOS)。 | | **存储与 Broker** | **Log Segment** | Partition 内的日志文件片段。 | 一个 Partition 的日志由多个 Log Segment 组成,方便日志滚动和清理。**类比 RocketMQ CommitLog 的切分文件**。 | | | **Log 结构** | Kafka 高性能的基石,Partition 是一个只追加的、有序的日志文件。 | 顺序写磁盘,避免了随机写的巨大开销,性能极高。 | | | **零拷贝** | 消费者高性能的关键,数据在内核态直接从磁盘发送到网卡。 | 避免了数据在用户空间和内核空间之间的多次拷贝,极大提升了吞吐量。 | | | **索引** | 加速消息查找,包括 Offset 索引和时间戳索引。 | 允许消费者根据 Offset 或时间戳快速定位消息,而不必遍历整个日志。 | | | **Retention Policy** | 消息保留策略,按时间或大小删除旧消息。 | Kafka 的清理机制,防止磁盘被占满,如保留7天或达到10GB大小。 | | **消费者端** | **Rebalance** | 消费者组内分区分配的动态过程,是负载均衡和故障转移的基础。 | 当组内消费者数量或 Topic 分区数量变化时触发,期间消费会短暂停止。 | | | **提交机制** | 决定消费语义的关键,消费者如何告知 Broker 自己的消费进度。 | 分为自动提交、手动同步提交(可靠但慢)、手动异步提交(快但需处理回调)。 | | | **消费语义** | 消息传递的保证级别。 | At-Most-Once (最多一次), At-Least-Once (至少一次,默认), Exactly-Once (精确一次)。 | | | **位移主题** | 消费位移的存储位置。 | Kafka 将消费组的 Offset 信息存储在内部的 Topic `**__consumer_offsets**` 中,而非 Broker 文件。 | | **集群协调与高可用** | **ISR (In-Sync Replica)** | 副本同步集合,确保消息可靠复制。 | 与 Leader 保持同步的 Follower 集合。acks=all 时,Leader 必须等待 ISR 中所有副本确认。**类比 RocketMQ 的 Master-Slave 同步机制**。 | | | **KRaft 模式** | Kafka 的未来方向,移除 Zookeeper 依赖。 | 使用 Raft 协议在 Broker 内部选举 Controller,简化部署和运维,提升性能。 | | | **Cooperative Rebalance** | 更平滑的 Rebalance 协议,解决 Rebalance 时的“Stop-the-World”问题。 | 分区可以增量地、分批地重新分配,消费中断更短,体验更平滑。 | | **高级特性** | **消息压缩** | 节省网络带宽和磁盘空间。 | Producer 可在发送前对消息进行压缩(如 Gzip, Snappy),Broker 和 Consumer 自动解压。 | | | **安全机制** | 保障数据安全。 | 包括 SSL/TLS (加密传输)、SASL (身份认证)、ACL (权限控制)。 | **kafka VS RocketMQ 能力对比** | 特性 | Kafka | RocketMQ | | ---------------- | ---------------------------------- | ---------------------------------- | | 系统类型 | 分布式消息中间件 | 分布式消息中间件 | | 核心功能 | 异步解耦、削峰填谷、异步通信 | 异步解耦、削峰填谷、事务消息 | | 典型角色 | Producer / Broker / Consumer | Producer / Broker / Consumer | | 数据模型 | Topic → Partition → Offset | Topic → MessageQueue → Offset | | 消费模式 | 基于 **Pull(拉取)** | 基于 **Pull(主动拉)** | | 消息投递 | “至少一次” (At least once) | “至少一次”,支持事务一致性 | | 支持事务 | 有(Kafka Transaction) | 有(RocketMQ Transaction Message) | | 是否支持本地事务 | ❌ 否。Kafka 事务只管消息投递原子性 | ✅ 是。业务代码可以和消息发送绑定 | | 存储方式 | Partition 文件(日志结构) | CommitLog + ConsumeQueue | | 持久化机制 | Append-only + LogSegment | CommitLog 顺序写 | | Offset 管理 | Consumer 自管(可提交) | Broker 统一维护(集群模式) | | 路由发现 | Controller(或 ZooKeeper) | NameServer | | 高可用 | ISR 同步副本机制 | 主从 + 同步复制机制 | | 消息清理 | TTL 或磁盘大小策略 | 定期清理 CommitLog | | 应用方向 | 大数据、日志流处理、实时分析 | 分布式系统、事务一致性、业务异步化 | 本地事务: 不是消息队列的事务,而是 **你自己服务内部** 的那部分业务逻辑(mysql,Spring事务等)。 RocketMQ 可以做到 让 **本地事务** 和 **消息发送** 的结果保持一致。 Kafka做不到。 ## 三个消息队列对比 | 特性维度 | RabbitMQ | RocketMQ | Kafka | | ---------------- | ---------------------------------------------- | ------------------------------------------------ | ---------------------------------------------- | | **架构模型** | **中心化 Broker** (AMQP 协议) | **分布式消息平台** | **分布式日志系统** | | **核心设计目标** | **智能消息路由**和**可靠传递** | **金融级可靠性**和**海量消息**的强一致性 | **高吞吐量**的**实时数据流处理** | | **消息模型** | **Exchange -> Queue** 模型,路由灵活 | **Topic -> MessageQueue** 模型,借鉴 Kafka | **Topic -> Partition** 模型,简单直接 | | **性能吞吐** | 中等(万级/秒) | **高**(十万级/秒) | **极高**(百万级/秒) | | **可靠性** | 高(通过 ACK、持久化、镜像队列) | **极高**(同步/异步刷盘、事务消息、主从复制) | 高(通过多副本 ISR 机制) | | **顺序性** | **队列内有序**(多消费者会打乱) | **队列内有序**(严格顺序消息) | **分区内有序**(全局无序) | | **消息回溯** | 不支持(需插件) | **原生支持**(按 Offset 或时间) | **原生支持**(按 Offset 或时间) | | **事务消息** | **不支持**(可通过 AMQP 事务模拟,性能差) | **原生支持**(两阶段提交+状态回查) | **不支持**(需借助外部系统,如 Kafka Streams) | | **生态系统** | 成熟,支持多种协议 | 主要在 Java 生态,阿里系生态完善 | **极其庞大**(流处理、大数据生态) | | **运维复杂度** | **较低**(集群搭建相对简单) | **中等**(Java 技术栈,对国内团队友好) | **较高**(依赖 Zookeeper,配置复杂) | | **典型应用场景** | **微服务通信**、**业务事件驱动**、**任务队列** | **订单交易**、**金融支付**、**大规模互联网应用** | **日志收集**、**大数据管道**、**实时监控** | 哪个消息队列应用最广泛? 从全球范围来看,**Kafka 和 RabbitMQ 的应用广度远超 RocketMQ**。RocketMQ在国内应用广泛,毕竟起源于国内,专为 Java 生态而生。 1. **Kafka:大数据领域的绝对霸主** - - 在**大数据、日志聚合、流式计算**领域,Kafka 是事实上的标准。几乎所有的大数据组件(Flink, Spark, Storm)都与 Kafka 深度集成。 - 如果你的场景是处理海量的数据流,Kafka 几乎是首选。**在这个领域,它的应用最广泛。** 1. **RabbitMQ:企业级应用和微服务的常青树** - - 在**传统企业应用、微服务架构、业务系统集成**方面,RabbitMQ 凭借其成熟稳定、灵活的路由和丰富的协议支持,拥有巨大的用户基数。 - 很多公司的第一个消息队列就是 RabbitMQ,因为它**上手快**,能解决大部分业务解耦和异步通信问题。在通用业务集成领域,它的应用非常广泛。 1. **RocketMQ:亚洲和大型互联网公司的利器** - - RocketMQ 起源于阿里,在国内的**大型互联网公司、金融、电商**领域应用极为广泛。它经历了“双十一”等超大规模场景的考验。 - 它的优势在于对**金融级事务消息、顺序消息、海量消息堆积**等复杂场景的完美支持。 - 不过,在欧美市场,RocketMQ 的知名度和采用率远不及 Kafka 和 RabbitMQ。

6/12-6/13 分布式消息队列(1)

### 1、消息队列优势 ​ **异步**,生产者发布后无需等待消费者,可以继续执行自己的的逻辑。 ​ **削峰**,将大量的请求保存到队列中,服务器依据自身能力逐步处理请求。 ### 2、分布式消息队列优势 ​ **数据持久化**,可以把消息存储到硬盘,数据不会丢失。 ​ **可扩展性**,可根据需求动态增加或者减少结点,保持服务的稳定 ​ **应用解耦**,允许不同框架语言的系统进行数据的传输和读取 (**解耦**:生产者和消费者之间不直接依赖对方的实现细节(如接口、调用方式、运行状态等),双方只需约定数据格式或协议) ### 3、rabbitmq基本使用 #### 3.1 helloworld ​ 生产者 注意: ​ 1)远程连接时要打开5672端口网络安全组,也要打开服务器本地防火墙 ​ 2)channel可以复用connection的tcp连接,且每个channel也可以有自己的状态。就像一家快递公司(Connection)可以派多个快递员(Channel)同时送货,比每次送货都重新开一家快递公司高效多了。 ```java import com.rabbitmq.client.Channel; import com.rabbitmq.client.Connection; import com.rabbitmq.client.ConnectionFactory; import java.nio.charset.StandardCharsets; public class Send { private final static String QUEUE_NAME = "hello"; private final static String REMOTE_HOST = "x.x.x.x";// 输入你的服务器 public static void main(String[] argv) throws Exception { // 创建连接工厂 ConnectionFactory factory = new ConnectionFactory(); factory.setHost(REMOTE_HOST); factory.setUsername("klee"); // 账号 factory.setPassword(""); // 密码 // 创建连接 try (Connection connection = factory.newConnection(); Channel channel = connection.createChannel()) { /** * 队列名 * 是否持久化 * 是否独占发送连接 * 是否自动删除队列 * 其他参数 */ channel.queueDeclare(QUEUE_NAME, false, false, false, null); String message = "Hello World2!"; // 发送消息 channel.basicPublish("", QUEUE_NAME, null, message.getBytes(StandardCharsets.UTF_8)); System.out.println(" [x] Sent '" + message + "'"); } } } ``` ​ 消费者 ```java import com.rabbitmq.client.Channel; import com.rabbitmq.client.Connection; import com.rabbitmq.client.ConnectionFactory; import com.rabbitmq.client.DeliverCallback; import java.nio.charset.StandardCharsets; public class Recv { private final static String QUEUE_NAME = "hello"; private final static String REMOTE_HOST = ""; public static void main(String[] argv) throws Exception { // 创建连接工厂 ConnectionFactory factory = new ConnectionFactory(); factory.setHost(REMOTE_HOST); factory.setUsername("klee"); factory.setPassword(""); // 从工厂获取一个新的连接 Connection connection = factory.newConnection(); // 从连接中创建一个新的频道 Channel channel = connection.createChannel(); // 创建队列,在该频道上声明我们正在监听的队列 channel.queueDeclare(QUEUE_NAME, false, false, false, null); // 在控制台打印等待接收消息的信息 System.out.println(" [*] Waiting for messages. To exit press CTRL+C"); // 定义了如何处理消息,创建一个新的DeliverCallback来处理接收到的消息 DeliverCallback deliverCallback = (consumerTag, delivery) -> { // 将消息体转换为字符串 String message = new String(delivery.getBody(), StandardCharsets.UTF_8); // 在控制台打印已接收消息的信息 System.out.println(" [x] Received '" + message + "'"); }; // 在频道上开始消费队列中的消息,接收到的消息会传递给deliverCallback来处理,会持续阻塞 channel.basicConsume(QUEUE_NAME, true, deliverCallback, consumerTag -> { }); } } ``` #### 3.2 workqueue ​ 生产者 ```java import com.rabbitmq.client.Channel; import com.rabbitmq.client.Connection; import com.rabbitmq.client.ConnectionFactory; import com.rabbitmq.client.MessageProperties; import java.util.Scanner; public class NewTask { // 定义要使用的队列名称,这次的队列就改叫multi_queue private static final String TASK_QUEUE_NAME = "task_queue"; public static void main(String[] argv) throws Exception { // 创建一个连接工厂 ConnectionFactory factory = new ConnectionFactory(); // 设置RabbitMQ服务的主机名 factory.setHost(""); factory.setUsername("klee"); factory.setPassword(""); // 创建一个新的连接 try (Connection connection = factory.newConnection(); // 创建一个新的频道 Channel channel = connection.createChannel()) { // 声明队列参数,包括队列名称、是否持久化等 channel.queueDeclare(TASK_QUEUE_NAME, false, false, false, null); // 创建一个输入扫描器,用于读取控制台输入 Scanner scanner = new Scanner(System.in); // 使用循环,每当用户在控制台输入一行文本,就将其作为消息发送 while (scanner.hasNext()) { // 读取用户在控制台输入的下一行文本 String message = scanner.nextLine(); // 发布消息到队列,设置消息持久化 channel.basicPublish("", TASK_QUEUE_NAME, MessageProperties.PERSISTENT_TEXT_PLAIN, message.getBytes("UTF-8")); // 输出到控制台,表示消息已发送 System.out.println(" [x] Sent '" + message + "'"); } } } } ``` ​ 多消费者: 注意: ​ 1)channel.basicQos(X); 启动则开启 **公平** 轮巡,每个消费者最多积压 X 个未确认消息 ​ 2)channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false); 确认信息,慎用批量确认 ​ 3)channel.basicNack(delivery.getEnvelope().getDeliveryTag(), false, false); ​ 最后一个参数true为重新入队,false为拒绝或进入死信队列,慎用true小心消息循环 ```java import com.rabbitmq.client.Channel; import com.rabbitmq.client.Connection; import com.rabbitmq.client.ConnectionFactory; import com.rabbitmq.client.DeliverCallback; import lombok.extern.slf4j.Slf4j; @Slf4j public class Worker { private static final String TASK_QUEUE_NAME = "task_queue"; public static void main(String[] argv) throws Exception { ConnectionFactory factory = new ConnectionFactory(); factory.setHost(""); factory.setUsername("klee"); factory.setPassword(""); final Connection connection = factory.newConnection(); for(int i=1;i<=2;i++) { final Channel channel = connection.createChannel(); // 声明队列 channel.queueDeclare(TASK_QUEUE_NAME, false, false, false, null); System.out.println("[*] 消费者" + i + "号 Waiting for messages. To exit press CTRL+C"); // 启动则开启公平轮巡,每个消费者最多积压 prefetchCount 个未确认消息 channel.basicQos(1); int finalI = i; // 接收消息 DeliverCallback deliverCallback = (consumerTag, delivery) -> { String message = new String(delivery.getBody(), "UTF-8"); System.out.println(" [x] Received '" + message + "'"); try { doWork(message, finalI); } catch (InterruptedException e) { log.error("doWork error", e); // 拒绝消息 channel.basicNack(delivery.getEnvelope().getDeliveryTag(), false, false); } finally { System.out.println(" [x]" + message + " Done by 消费者" + finalI); // 确认消息已经被处理 channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false); } }; // 开启消息接收,建议手动确认 channel.basicConsume(TASK_QUEUE_NAME, false, deliverCallback, consumerTag -> { }); } } private static void doWork(String task,int id) throws InterruptedException { System.out.println("任务:" + task + "已经被消费者id:" + id + "接收"); Thread.sleep(30*1000); // 模拟处理任务需要时间 } } ``` #### 3.3 fanout交换机 ​ 即全体广播,**生产者**声明一个交换机并指定其类型为fanout,发布消息也指定交换机即可 ```java package com.zxw.bi_intelligence.rabbitmq.exchanger.fanout; import com.rabbitmq.client.*; import java.util.Scanner; public class EmitLog { private static final String EXCHANGE_NAME = "fanout"; public static void main(String[] argv) throws Exception { ConnectionFactory factory = new ConnectionFactory(); // 设置RabbitMQ服务的主机名 factory.setHost(""); factory.setUsername("klee"); factory.setPassword(""); try (Connection connection = factory.newConnection(); Channel channel = connection.createChannel()) { // 创建一个交换机 channel.exchangeDeclare(EXCHANGE_NAME, BuiltinExchangeType.FANOUT); // 创建一个输入扫描器,用于读取控制台输入 Scanner scanner = new Scanner(System.in); // 使用循环,每当用户在控制台输入一行文本,就将其作为消息发送 while (scanner.hasNext()) { // 读取用户在控制台输入的下一行文本 String message = scanner.nextLine(); // 发布消息到交换机 channel.basicPublish(EXCHANGE_NAME, "", null, message.getBytes("UTF-8")); // 输出到控制台,表示消息已发送 System.out.println(" [x] Sent '" + message + "'"); } } } } ``` ​ **消费者** ​ 要注意把队列和交换机绑定即可, channel.queueBind(队列名称, 交换机名称, ""); ```java import com.rabbitmq.client.*; public class ReceiveLogs { private static final String EXCHANGE_NAME = "fanout"; private static final String FANOUT_QUEUE_ONE = "logsRepo"; private static final String FANOUT_QUEUE_TWO = "mirrorRepo"; public static void main(String[] argv) throws Exception { ConnectionFactory factory = new ConnectionFactory(); factory.setHost(""); factory.setUsername("klee"); factory.setPassword(""); Connection connection = factory.newConnection(); // 创建两个队列,一个专们给日志,另一个备份 Channel channel = connection.createChannel(); channel.exchangeDeclare(EXCHANGE_NAME, BuiltinExchangeType.FANOUT); channel.queueDeclare(FANOUT_QUEUE_ONE, false, false, false, null); channel.queueBind(FANOUT_QUEUE_ONE, EXCHANGE_NAME, ""); System.out.println(" [*] Waiting for messages. To exit press CTRL+C"); DeliverCallback deliverCallback = (consumerTag, delivery) -> { String message = new String(delivery.getBody(), "UTF-8"); System.out.println(" [x] Received 消费者1取出:"+ message + "'"); }; channel.basicConsume(FANOUT_QUEUE_ONE, true, deliverCallback, consumerTag -> { }); Channel channel2 = connection.createChannel(); channel2.exchangeDeclare(EXCHANGE_NAME, BuiltinExchangeType.FANOUT); channel2.queueDeclare(FANOUT_QUEUE_TWO, false, false, false, null); channel2.queueBind(FANOUT_QUEUE_TWO, EXCHANGE_NAME, ""); System.out.println(" [*] Waiting for messages. To exit press CTRL+C"); DeliverCallback deliverCallback2 = (consumerTag, delivery) -> { String message = new String(delivery.getBody(), "UTF-8"); System.out.println(" [x] Received 消费者2取出:"+ message + "'"); }; channel.basicConsume(FANOUT_QUEUE_TWO, true, deliverCallback2, consumerTag -> { }); } } ``` #### 3.4 direct交换机 ​ 特点:按照交换机上的路由表转发给指定的队列 ​ 生产者: ​ 声明交换机和类型channel.exchangeDeclare(EXCHANGE_NAME, BuiltinExchangeType.DIRECT); ```java import com.rabbitmq.client.BuiltinExchangeType; import com.rabbitmq.client.Channel; import com.rabbitmq.client.Connection; import com.rabbitmq.client.ConnectionFactory; import java.util.Scanner; public class EmitLogDirect { private static final String EXCHANGE_NAME = "direct"; public static void main(String[] argv) throws Exception { ConnectionFactory factory = new ConnectionFactory(); factory.setHost(""); factory.setUsername(""); factory.setPassword(""); try (Connection connection = factory.newConnection(); Channel channel = connection.createChannel()) { // channel.exchangeDeclare(EXCHANGE_NAME, BuiltinExchangeType.DIRECT); // 创建一个输入扫描器,用于读取控制台输入 Scanner scanner = new Scanner(System.in); // 使用循环,每当用户在控制台输入一行文本,就将其作为消息发送 while (scanner.hasNext()) { // 读取用户在控制台输入的下一行文本 String raw = scanner.nextLine(); // 如果用户输入的文本行包含两个单词,就将其作为消息和路由键发送 String[] s = raw.split(" "); if(s.length != 2){ continue; } String message = s[0]; String routerKey = s[1]; // 发布消息到交换机 channel.basicPublish(EXCHANGE_NAME, routerKey, null, message.getBytes("UTF-8")); // 输出到控制台,表示消息已发送 System.out.println(" [x] Sent " + message + " by direct_exchanger"); } } } } ``` ​ 消费者: ​ 公式:1)交换机声明 ​ 2)队列声明 ​ 3)交换机、队列绑定,此时要声明routerKey以便于交换机定向转发 ```java import com.rabbitmq.client.*; public class ReceiveLogsDirect { private static final String EXCHANGE_NAME = "direct"; public static void main(String[] argv) throws Exception { ConnectionFactory factory = new ConnectionFactory(); factory.setHost(""); factory.setUsername("klee"); factory.setPassword(""); Connection connection = factory.newConnection(); // 声明两个队列,一个员工小郑任务队列,一个员工鱼皮任务队列 Channel channel = connection.createChannel(); // 声明交换机,并指定其类型为direct channel.exchangeDeclare(EXCHANGE_NAME, BuiltinExchangeType.DIRECT); String queueName = "小郑的任务单"; channel.queueDeclare(queueName, false, false, false, null); // 绑定队列到交换机,指定routingKey channel.queueBind(queueName, EXCHANGE_NAME, "Xzheng"); System.out.println(" [ 小郑 ] Waiting for messages. To exit press CTRL+C"); DeliverCallback deliverCallback = (consumerTag, delivery) -> { String message = new String(delivery.getBody(), "UTF-8"); System.out.println(" [ 小郑 ] Received '" + delivery.getEnvelope().getRoutingKey() + "':'" + message + "'"); }; channel.basicConsume(queueName, true, deliverCallback, consumerTag -> { }); Channel c2 = connection.createChannel(); c2.exchangeDeclare(EXCHANGE_NAME, BuiltinExchangeType.DIRECT); String qN2 = "鱼皮的任务单"; c2.queueDeclare(qN2, false, false, false, null); // 绑定队列到交换机,指定routingKey c2.queueBind(qN2, EXCHANGE_NAME, "Yupi"); System.out.println(" [ 鱼皮 ] Waiting for messages. To exit press CTRL+C"); DeliverCallback deliverCallback2 = (consumerTag, delivery) -> { String message = new String(delivery.getBody(), "UTF-8"); System.out.println(" [ 鱼皮 ] Received '" + delivery.getEnvelope().getRoutingKey() + "':'" + message + "'"); }; c2.basicConsume(qN2, true, deliverCallback2, consumerTag -> { }); } } ```

阿里Linux3(基于CentOS8)部署rabbitmq

## 下载Rabbitmq 在Github官网上可以找到所有版本的Rabbitmq-server,挑选一个稳定的版本下载 https://github.com/rabbitmq/rabbitmq-server/tags ## 下载erlang Erlang和RabbitMQ版本对照:https://www.rabbitmq.com/which-erlang.html 这里我安装的rabbitmq是3.13版本的,因此我选的erlang是26大版本,各位鱼友挑选对应的版本 erlang下载地址:https://packagecloud.io/rabbitmq/erlang 其中el8代表——CentOs8系统 <img src="https://pic.code-nav.cn/post_picture/1830958011788279809/pH6WtAyvCYh79AVw.webp" alt="image.png" width="100%" /> 下载完成后将两个rpm文件都放到linux服务器上 <img src="https://pic.code-nav.cn/post_picture/1830958011788279809/jtnZMvRsAvg5z4JO.webp" alt="image.png" width="100%" /> 我放置在了/usr/rabbitmq的目录中 ## 解压安装 到命令行中执行以下命令:(部分命令失败可能因为权限不足) ``` # 解压 rpm -Uvh erlang-26.2.5.3-1.el8.x86_64.rpm # 安装 yum install -y erlang # 检验erl是否安装成功 erl -v ``` ### 安装Rabbitmq Rabbitmq的安装需要该插件 ``` yum install -y socat ``` 安装Rabbitmq ``` # 解压 rpm -Uvh rabbitmq-server-3.13.7-1.el8.noarch.rpm # 安装 yum install -y rabbitmq-server ``` ## 启动rabbitmq ``` # 启动rabbitmq systemctl start rabbitmq-server # 查看rabbitmq状态 systemctl status rabbitmq-server ``` 看到这里是active就成功运行起来了 <img src="https://pic.code-nav.cn/post_picture/1830958011788279809/zIyafrAxZGAAhUTT.webp" alt="image.png" width="100%" /> 默认情况下rabbitmq没有安装web端的客户端软件,需要我们自己安装 ## 打开RabbitMQWeb管理界面插件 ```rabbitmq-plugins enable rabbitmq_management``` ## 远程连接 目前我们做的仅仅是将Rabbitmq在服务器上运行起来,但是我们的电脑上仍然无法访问直接通过ip访问 ### 打开防火墙 **这里要注意!!5672是rabbitmq的端口,15672是rabbitmq管理页面的端口** 宝塔linux的防火墙: 服务器的防火墙: <img src="https://pic.code-nav.cn/post_picture/1830958011788279809/eKZU6i1PetpmoXZ9.webp" alt="image.png" width="100%" /> 配置完成后我们暂时可以通过```ip:15672```访问到rabbitmq管理界面的登录页 尝试用``` username:guest password:guest ```登录 **你会发现登陆不进去!!!** 这是因为这样的登录方式仅允许localhost登录,因此我们是无法直接通过默认账密进行登录,必须自己添加一个远程登录的用户 ## 创建远程登录用户 ``` # 添加用户 rabbitmqctl add_user 用户名 密码 # 设置用户角色,分配操作权限 rabbitmqctl set_user_tags 用户名 角色 # 为用户添加资源权限(授予访问虚拟机根节点的所有权限) rabbitmqctl set_permissions -p / 用户名 ".*" ".*" ".*" ``` **角色有四种**: - `administrator`:可以登录控制台、查看所有信息、并对rabbitmq进行管理 - `monToring`:监控者;登录控制台,查看所有信息 - `policymaker`:策略制定者;登录控制台指定策略 - `managment`:普通管理员;登录控制 ``` # 检验是否成功创建用户并分配资源 rabbitmqctl list_permissions -p / ``` 如果能成功显示则完成配置 <img src="https://pic.code-nav.cn/post_picture/1830958011788279809/5jxZQq7SUY5IT6WB.png" alt="image.png" width="665px" /> 这时候,你就可以在浏览器上通过公网访问rabbitmq的管理界面了 <img src="https://pic.code-nav.cn/post_picture/1830958011788279809/7n7RvjvUsI9XHR4G.webp" alt="image.png" width="100%" /> ## JAVA连接Rabbitmq 导入相关依赖: ``` <!-- https://mvnrepository.com/artifact/com.rabbitmq/amqp-client --> <dependency> <groupId>com.rabbitmq</groupId> <artifactId>amqp-client</artifactId> <version>5.21.0</version> </dependency> ``` 修改配置文件: ``` spring: rabbitmq: host: xx.xx.xx.xx username: jerry_chen password: xxxxxx ``` 使用客户端正常连接,传入对应的host,username,password ``` ConnectionFactory factory = new ConnectionFactory(); factory.setHost(host); factory.setUsername(username); factory.setPassword(password); Connection connection = factory.newConnection(); Channel channel = connection.createChannel(); ``` 成功通过Java连接到服务器的rabbitmq <img src="https://pic.code-nav.cn/post_picture/1830958011788279809/DNrxBI3u1GVILVpm.webp" alt="image.png" width="100%" />

#学习笔记# #Java# #rabbitmq# 因为要做俱乐部的项目,去学习了Rabbitmq,尚硅谷的课程总的来说还是可以的,但是后面的负载均衡(p86之后)存在课程缺失,没有详细讲解和操作,直接带过,建议小伙伴们可以去看一下其他的教程 rabbitmq是当前最流行的消息中间件之一,他是在AMQP(高级消息队列协议)基础上完成的,是使用Erlang语言开发的,具有高并发的特性。MQ可以比喻成一个消息中转站,不处理消息,只是进行存储和转发,在鱼总的分享项目的直播--尚医通中也有介绍 最后分享一波笔记,安装教程,centos7安装教程(韩顺平)

下载 APP