# 仿RabbitMQ实现消息队列 **Repository Path**: Faiz--555/MQ ## Basic Information - **Project Name**: 仿RabbitMQ实现消息队列 - **Description**: 在实际的后端开发中,尤其是在分布式系统里,跨主机之间使用生产者消费者模型是非常普遍的需求。因此,我们通常会把阻塞队列封装成一个独立的服务器程序,并且赋予其更丰富的功能。这样的服务程序我们就称为消息队列(Message Queue, MQ)。 其中 RabbitMQ 是一个非常知名、功能强大且广泛使用的消息队列。本项目就是仿照 RabbitMQ 模拟实现一个简单的消息队列。 - **Primary Language**: C++ - **License**: Not specified - **Default Branch**: master - **Homepage**: None - **GVP Project**: No ## Statistics - **Stars**: 2 - **Forks**: 0 - **Created**: 2024-07-15 - **Last Updated**: 2025-08-18 ## Categories & Tags **Categories**: Uncategorized **Tags**: None ## README # 仿RabbitMQ实现消息队列 —— 整体框架 ### 博客地址 [仿RabbitMQ实现消息队列](https://blog.csdn.net/decade777555/article/details/140757156?spm=1001.2014.3001.5501) ## 一、项目简介 在实际的后端开发中,尤其是在分布式系统里,跨主机之间使用生产者消费者模型是非常普遍的需求。因此,我们通常会把阻塞队列封装成一个独立的服务器程序,并且赋予其更丰富的功能。这样的服务程序我们就称为消息队列(Message Queue, MQ)。 其中 RabbitMQ 是一个非常知名、功能强大且广泛使用的消息队列。本项目就是仿照 RabbitMQ 模拟实现一个简单的消息队列。 ### 需求分析 - **生产者 (Producer)** - **消费者 (Consumer)** - **中间人 (Broker)** - **发布 (Publish)** - **订阅 (Subscribe)** 我们需要实现的内容包括: 1. Broker 服务器:消息队列代理服务器 2. 消息发布客户端:向服务器发布消息 (生产者) 3. 消息订阅客户端:从服务器订阅消息 (消费者) 我们的消息队列是基于对 AMQP 协议的理解进行的整合,那么什么是 AMQP 协议呢: AMQP (Advanced Message Queuing Protocol) 是一种网络协议,主要用于消息中间件之间的异步通信。AMQP 提供了一种标准的方式来发送和接收消息,使得不同厂商的消息中间件可以互相操作。 ### AMQP 特点 - **开放标准**:AMQP 是一个开放标准,这意味着它的规范是公开的,并且任何人都可以实现它。 - **二进制协议**:AMQP 使用二进制编码,这使得它比基于文本的协议更高效。 - **可靠性**:AMQP 设计为确保消息的可靠传输,包括确认机制、事务支持等。 - **灵活性**:AMQP 支持多种消息路由模式,如点对点(point-to-point) 和发布/订阅(publish/subscribe)。 - **互操作性**:AMQP 允许不同厂商的消息中间件互相通信,不受客户端或中间件使用的编程语言的影响。 ### AMQP 模型 - **Broker (消息代理)**:这是消息中间件的核心组件,负责接收、存储和转发消息。 - **Exchange (交换器)**:Exchange 接收来自生产者的消息,并根据配置规则将消息发送到一个或多个队列。 - **Queue (队列)**:队列是用来存储消息的数据结构,直到消费者取走这些消息。 - **Binding (绑定)**:绑定定义了 Exchange 和 Queue 之间的关系,确定消息如何从 Exchange 到达 Queue。 - **Producer (生产者)**:生产者是向消息中间件发送消息的应用程序。 - **Consumer (消费者)**:消费者是从消息中间件接收消息的应用程序。 ## 二、服务端模块 服务端模块包括: 1. **交换机数据管理** 2. **队列数据管理** 3. **绑定数据管理** 4. **消息数据管理** 5. **虚拟机数据管理** 6. **路由匹配管理** 7. **消费者管理** 8. **信道管理** 9. **连接管理** 10. **服务器模块** - 生产者消费者则通过网络发送请求来调用这些API,实现生产者消费者模型 - 交换机类型 - 本项目实现三种交换机类型,也是最常见的: • Direct: ⽣产者发送消息时, 直接指定被该交换机绑定的队列名 • Fanout: ⽣产者发送的消息会被复制到该交换机的所有队列中 • Topic: 绑定队列到交换机上时, 指定⼀个字符串为 bindingKey。发送消息指定⼀个字符串 routingKey。当 routingKey 和 bindingKey 满⾜⼀定的匹配条件的时候, 则 把消息投递到指定队列 持久化 交换机,队列,绑定,消息都是需要持久化的,我们需要根据持久化来保证在程序或主机重启时数据不会丢失。在项目中我们使用了Sqlite数据库来进行本地的轻量级存储 网络通信 ⽣产者和消费者都是客⼾端程序, Broker 则是作为服务器,通过⽹络进⾏通信。 我们在broker的基础上,再加上建立连接和打开信道的操作,这样 可以更好地复用TCP连接,达到长连接的效果,避免频繁的创建关闭TCP连接。 ### 1、交换机数据管理 交换机数据管理就是描述了交换机应该有哪些数据 我们可以设置交换机的类型以及消息的基本属性,基本结构如下: ``` syntax ='proto3'; package bitmq; enum ExchangeType { UNKNOWTYPE=0; DIRECT=1; FANOUT=2; TOPIC=3; }; enum DeliverMode { UNKNOWMODE=0; UNDURABLE=1; DURABLE=2; }; message BasicProperties { string id=1; DeliverMode delivery_mode=2; string routing_key=3; }; message message { message Payload { BasicProperties properties=1; string body=2; string valid=3; }; Payload payload=1; uint32 offset=2; uint32 length=3; }; ``` - 1、交换机的名称:也是交换机的唯一标识 - 2、交换机的类型:决定了消息的转发方式(三种 ) 每个队列与交换机绑定信息中有binding_key,每条消息中有routing_key 1.直接交换:binding_key与routing_key相同时,将消息放入队列 2.广播交换:交换机绑定的所有队列都放入消息 3.主题交换:根据具体的匹配算法,将符合匹配条件的binding_key和routing_key对应的队列放入消息。 - 3、持久化标志:决定当前交换机的数据是否需要持久化存储 - 4、自动删除标志:如果关联该交换机的客户端都退出了,是否需要自动删除交换机。 - 5、交换机的其他参数,在本项目中未具体使用。 ### 2、队列数据管理 队列数据管理涉及队列的基本操作及信息管理,包括但不限于: - 创建队列 - 销毁队列 - 获取指定队列信息 - 获取队列数量 - 获取所有队列的名称 ### 3、绑定数据管理 绑定数据管理描述队列与交换机之间的关系,主要功能包括: - 添加绑定信息 - 解除绑定信息 - 获取交换机所有的相关绑定信息 - 获取队列所有的绑定信息 - 获取绑定信息的数量 ### 4、消息数据管理 消息数据管理包括消息的存储及操作: - 向队列新增消息 - 获取队首消息 - 对消息进行确认 - 恢复队列的历史消息 - 垃圾回收 ### 5、虚拟机数据管理 虚拟机数据管理涉及虚拟环境下的资源分配和管理: - 管理虚拟机内的交换机、队列、绑定、消息等资源 - 虚拟机作为一个逻辑上的容器,用于隔离不同的应用环境 ### 6、路由匹配管理 路由匹配管理负责根据消息的属性将消息路由到正确的队列: - 根据交换机类型和消息的routing key匹配队列 - 实现Direct、Fanout、Topic等交换机类型的消息匹配逻辑 ### 7、消费者管理 消费者管理负责维护消费者的状态和行为: - 增加消费者 - 删除消费者 - 获取指定队列的消费者 - 获取队列的消费者列表是否为空 - 管理消费者信息,例如消费者标识、订阅的队列名称等 ### 8、信道管理 信道管理提供了通信的逻辑通道: - 打开一个信道 - 关闭一个信道 - 获取指定信道句柄 - 管理信道ID、关联的消费者句柄、虚拟机句柄以及工作线程池句柄 ### 9、连接管理 连接管理负责网络连接的生命周期: - 新增连接 - 关闭连接 - 获取指定的连接信息 - 封装底层网络库中的连接概念,提供高级别的连接管理和操作接口 ### 10. 服务器模块 Broker 服务器模块是一个功能的整合,本质上这个模块并不提供实质的功能性操作,这个模块最重要的是资源的整合,是一个资源的载体。 - 一个服务器有一个工作线程池,其他所有的信道操作都是这一个线程池的。 - 一个服务器有一个虚拟机,其他所有的交换机,队列,绑定,消息的操作都是针对这个虚拟机进行的。 - 一个服务器有一个消费者管理。 - 通信相关的连接管理,协议处理模块句柄,也是一整个服务器有一套。 ## 三、客户端模块 客户端模块包括: 1. **消费者管理** 2. **信道管理** 3. **连接管理** 4. **异步线程池模块** ### 1. 消费者管理 消费者信息: - 消费者标识 - 订阅的队列名称 - 自动确认标志 - 消费处理回调函数指针 消费者管理:增删查操作 ### 2. 信道管理 所有提供的操作与服务端基本对应,因为客户端需要给用户提供什么服务,服务器就要给客户端提供什么服务。 - 信道提供的服务 - 声明/删除交换机 - 声明/删除队列 - 绑定信息的绑定/删除 - 消息的发布/订阅队列消息/取消订阅/队列消息的 Ack - 创建/关闭信道 ### 3. 连接管理 客户端连接的管理,本质是对客户端 TcpClient 的二次封装和管理。 对于用户,不需要有客户端的概念,连接对于用户来说就是客户端,通过连接创建信道来完成所需的服务,客户端这边的连接对用户来说就是一个资源的载体。 - 管理操作 - 连接服务器 - 创建信道 - 关闭信道 - 关闭连接 ### 4. 异步线程池模块 TcpClient 需要一个 EventLoopThread 模块进行 IO 事件监控。 收到推送的消息,需要对推送过来的消息进行处理,因此需要一个线程池来帮助我们完成消息处理的过程。 ![模块关系图](%E9%A1%B9%E7%9B%AE%E6%A8%A1%E5%9D%97%E5%85%B3%E7%B3%BB%E5%9B%BE.png)