消息队MQ

简介: 消息队MQ

文章描述

😊 @ 作者:Lion J
`💖 @ 主页:
https://blog.csdn.net/weixin_69252724`
🎉 @ 主题: 消息队列MQ_rabbitMQ搭建
⏱️ @ 创作时间:2024年03月9日
————————————————


举一个 电商的例子
在开发的一个场景中,用户下订单给订单服务,订单服务调用库存服务减产库存情况, 订单服务再下订单, 下完订单再通知用户订单信息

在这里插入图片描述


一、MQ是什么?

MQ 全称(Message Queue)又名消息队列,是一种提供消息队列服务的中间件,也称为消息中间件,是一套提供了消息生产、存储、消费全过程 API的软件系统(消息即数据)。通俗点说,就是一个先进先出的数据结构。

言简意赅的说,就是 将服务中间某个行为步骤先抽取到一个容器,让容器去操作,不影响当前的服务

二、常见MQ中间件

ZeroMQ: 号称最快的消息队列系统,尤其针对大吞吐量的需求场景。扩展性好,
开发比较灵活,采用 C 语言 实现,实际上只是一个 socket 库的重新封装,如果做为消息队列使用,需要开发大量的代码。 ZeroMQ 仅提供非持久性的队列,也就是说如果 down 机,数据将会丢失。
RabbitMQ: 使用 erlang 语言开发,性能较好,适合于企业级的开发。但是不利于做二次开发和维护。
ActiveMQ: 历史悠久的 Apache 开源项目。已经在很多产品中得到应用,实现
了 JMS1.1 规范,可以和 springjms 轻松融合,实现了多种协议,支持持久化到数据库,对队列数较多的情况支持不好。
RocketMQ: 阿里巴巴的 MQ 中间件,由 java 语言开发,性能非常好,能够撑住双十一的大流量,而且使用起来 很简单。
Kafka: Kafka 是 Apache 下的一个子项目,是一个高性能跨语言分布式Publish/Subscribe 消息队列系统, 相对于 ActiveMQ 是一个非常轻量级的消息系统,除了性能非常好之外,还是一个工作良好的分布式系统。

三、RocletMQ环境搭建

rabbitMQ搭建

  1. 下载解压
    https://rocketmq.apache.org/download/
  2. 配置环境变量
ROCKETMQ_HOME=D:\ProgramFiles\rocketmq-4.9.3
NAMESRV_ADDR =127.0.0.1:9876
  1. 启动Name Server

进入到bin目录输入命令:
mqnamesrv.cmd

  1. 启动Broker

进入到 bin 目录输入命令:
mqbroker.cmd -n 127.0.0.1:9876 atuoCreateTopicEnable=true

控制台安装启动

  1. 解压

在这里插入图片描述

  1. 修改其 src/main/resources 中的 application.properties 配置文件

在这里插入图片描述

  1. 在解压目录 rocketmq-console 的 pom.xml 中添加如下 JAXB 依赖。
<dependency>
<groupId>javax.xml.bind</groupId>
<artifactId>jaxb-api</artifactId>
<version>2.3.0</version>
</dependency>

<dependency>
<groupId>com.sun.xml.bind</groupId>
<artifactId>jaxb-impl</artifactId>
<version>2.3.0</version>
</dependency>

<dependency>
<groupId>com.sun.xml.bind</groupId>
<artifactId>jaxb-core</artifactId>
<version>2.3.0</version>
</dependency>

<dependency>
<groupId>javax.activation</groupId>
<artifactId>activation</artifactId>
<version>1.1.1</version>
</dependency>
  1. 打包_命令行进入到 rocketmq-console

mvn clean package -Dmaven.test.skip=true

  1. 打包后,进入 target 目录

启动控制台 java -jar rocketmq-console-ng-1.0.0.jar

  1. 访问

http://127.0.0.1:6060
在这里插入图片描述

四、RocketMQ架构

在这里插入图片描述

其中Broker是RocketMQ的核心, 当Broker启动后, 就会向NameServer中注册自身消息, 然后Producer在NameServer中获取Broker的信息,然后向Broker发送投递消息; 消费者Consumer在NameServer中获取Broker消息之后就会从Broker中接收消息

NameServer,Broker,Producer,Consumer。
Broker(邮递员) Broker 是 RocketMQ 的核心,负责消息的接收,存储,投递等功能.
NameServer(邮局) 消息队列的协调者,Broker 向它注册路由信息,同时Producer 和 Consumer 向其获取路由信息
Producer(寄件人) 消息的生产者,需要从 NameServer 获取 Broker 信息,然后与 Broker 建立连接,向 Broker 发送消 息
Consumer(收件人) 消息的消费者,需要从 NameServer 获取 Broker 信息,然后与 Broker 建立连接,从 Broker 获取消息
Topic(地区) 用来区分不同类型的消息,发送和接收消息前都需要先创建Topic,针对 Topic 来发送和接收消息
Message Queue(邮件) 为了提高性能和吞吐量,引入了 Message Queue,一个 Topic 可以设置一个或多个 Message Queue,这样消息就可以并行往各个Message Queue 发送消息,消费者也可以并行的从多个 Message Queue 读取消息 Message Message 是消息的载体。
Producer Group 生产者组,简单来说就是多个发送同一类消息的生产者称之为一个生产者组。
Consumer Group 消费者组,消费同一类消息的多个 consumer 实例组成一个消费者组。

五、java消息发送和接收演示

消息发送者

        public class MQProducerTest {
            public static void main(String[] args) throws Exception {
//1. 创建消息生产者, 指定生产者所属的组名
                DefaultMQProducer producer = new DefaultMQProducer("myproducer-group");
//2. 指定 Nameserver 地址
                producer.setNamesrvAddr("192.168.109.131:9876");
//3. 启动生产者
                producer.start();
//4. 创建消息对象,指定主题、标签和消息体
                Message msg = new Message("myTopic", "myTag",
                        ("RocketMQ Message").getBytes());
//5. 发送消息
                SendResult sendResult = producer.send(msg, 10000);
                System.out.println(sendResult);
//6. 关闭生产者
                producer.shutdown();
            }
        }

消息接收

        public class MQConsumerTest {
            public static void main(String[] args) throws Exception {
//1. 创建消息消费者, 指定消费者所属的组名
                DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("myconsumergroup");
//2. 指定 Nameserver 地址
                consumer.setNamesrvAddr("192.168.109.131:9876");
//3. 指定消费者订阅的主题和标签
                consumer.subscribe("myTopic", "*");
//4. 设置回调函数,编写处理消息的方法
                consumer.registerMessageListener(new MessageListenerConcurrently() {
                    @Override
                    public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt>
                                                                            msgs,
                                                                    ConsumeConcurrentlyContext
                                                                            context) {
                        System.out.println("Receive New Messages: " + msgs);//返回消费状态
                        return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
                    }
                });
//5. 启动消息消费者
                consumer.start();
                System.out.println("Consumer Started.");
            }
        }

六、案例

在这里插入图片描述

订单微服务发送消息

  1. 添加rocketmq依赖
<!--rocketmq-->
<dependency>
<groupId>org.apache.rocketmq</groupId>
<artifactId>rocketmq-spring-boot-starter</artifactId>
<version>2.0.2</version>
</dependency>
<dependency>
<groupId>org.apache.rocketmq</groupId>
<artifactId>rocketmq-client</artifactId>
<version>4.4.0</version>
</dependency>
  1. 添加配置

在这里插入图片描述

  1. 编写测试代码
@Autowired
private RocketMQTemplate rocketMQTemplate;
rocketMQTemplate.convertAndSend("order-topic", order);

用户微服务接收消息

  1. 添加依赖
        <!--rocketmq-->
        <dependency>
            <groupId>org.apache.rocketmq</groupId>
            <artifactId>rocketmq-spring-boot-starter</artifactId>
            <version>2.0.2</version>
        </dependency>
        <dependency>
            <groupId>org.apache.rocketmq</groupId>
            <artifactId>rocketmq-client</artifactId>
            <version>4.4.0</version>
        </dependency>
  1. 修配置文件

在这里插入图片描述

  1. 编写消息接收服务
@Service
@RocketMQMessageListener(consumerGroup = "shop-user", topic = "order-topic")
public class SmsService implements RocketMQListener<Order> {
    @Override
    public void onMessage(Order order) {
        System.out.println("收到一个订单信息:"+ JSON.toJSONString(order)+",接下来发送短信");
    }
}
  1. 启动服务,执行下单操作,观看后台输出

在这里插入图片描述

七、发送不同类型消息

RocketMQ 提供三种方式来发送普通消息:可靠同步发送、可靠异步发送、单向发送

可靠同步发送: 同步发送是指消息发送方发出数据后,会在收到接收方发回响应之后才发下一个数据包的通讯方 式。 此种方式应用场景非常广泛,例如重要通知邮件、报名短信通知、营销短信系统等。
可靠异步发送: 异步发送是指发送方发出数据后,不等接收方发回响应,接着发送下个数据包的通讯方式。发送 方通过回调接口接收服务器响应,并对响应结果进行处理。 异步发送一般用于链路耗时较长,对 RT 响应时间较为敏感的业务场景,例如用户视频上传后通知 启动转码服务,转码完成后通知推送转码结果等。
单向发送: 单向发送是指发送方只负责发送消息,不等待服务器回应且没有回调函数触发,即只发送请求不 等待应答。 适用于某些耗时非常短,但对可靠性要求并不高的场景,例如日志收集。

● 同步消息

//同步消息
//参数一: topic
//参数二: 消息内容
        SendResult sendResult = rocketMQTemplate.syncSend("test-topic-1", "这是一
                条同步消息");
                System.out.println(sendResult);

●发送异步消息

//参数一: topic
//参数二: 消息内容
//参数三: 回调函数, 处理返回结果
rocketMQTemplate.asyncSend("test-topic-1","这是一条异步消息",new

    SendCallback() {
        @Override
        public void onSuccess (SendResult sendResult){
            System.out.println(sendResult);
        }
        @Override
        public void onException (Throwable throwable){
            System.out.println(throwable);
        }
    });
//让线程不要终止
Thread.sleep(30000000)

●单向消息

rocketMQTemplate.sendOneWay("test-topic-1", "这是一条单向消息");
相关实践学习
快速体验阿里云云消息队列RocketMQ版
本实验将带您快速体验使用云消息队列RocketMQ版Serverless系列实例进行获取接入点、创建Topic、创建订阅组、收发消息、查看消息轨迹和仪表盘。
消息队列 MNS 入门课程
1、消息队列MNS简介 本节课介绍消息队列的MNS的基础概念 2、消息队列MNS特性 本节课介绍消息队列的MNS的主要特性 3、MNS的最佳实践及场景应用 本节课介绍消息队列的MNS的最佳实践及场景应用案例 4、手把手系列:消息队列MNS实操讲 本节课介绍消息队列的MNS的实际操作演示 5、动手实验:基于MNS,0基础轻松构建 Web Client 本节课带您一起基于MNS,0基础轻松构建 Web Client
相关文章
|
8月前
|
人工智能 运维 Serverless
一杯咖啡成本搞定多模态微调:FC DevPod + Llama-Factory 极速实战
告别显存不足、环境配置难、成本高昂的微调困境!基于阿里云函数计算FC与Llama-Factory,5分钟搭建微调流水线,一键完成多模态模型的微调。
777 20
|
4月前
|
消息中间件 存储 负载均衡
【消息队列MQ】消息队列MQ——三大核心作用:解耦、异步、削峰填谷;两大核心模型:点对点、发布订阅(附《 消息队列MQ 面试核心考点问答清单》)
本文系统解析消息队列MQ的**三大核心作用**(解耦、异步、削峰填谷)与**两大核心模型**(点对点、发布订阅),贯通原理、价值、实现、选型及避坑实践,构建分布式系统中MQ的全链路知识体系。
|
4月前
|
IDE 开发工具 Swift
Xcode 26.4.1 (17E202) 发布 - Apple 平台 IDE
IDE for iOS/iPadOS/macOS/watchOS/tvOS/visonOS
606 0
|
9月前
|
机器学习/深度学习 人工智能 前端开发
终端里的 AI 编程助手:OpenCode 使用指南
OpenCode 是开源的终端 AI 编码助手,支持 Claude、GPT-4 等模型,可在命令行完成代码编写、Bug 修复、项目重构。提供原生终端界面和上下文感知能力,适合全栈开发者和终端用户使用。
59787 11
|
人工智能 自然语言处理 算法
AI时代,ETL真的不行了吗?
本文探讨了AI技术如何深度参与数据处理与分析,推动企业数据集成从传统ETL向“ETL for AI”转型。通过分析AI与ETL的协作关系,指出未来数据集成将实现高效处理、安全流转与智能价值挖掘,助力企业迈向数智化转型。
AI时代,ETL真的不行了吗?
|
消息中间件 存储 负载均衡
2024消息队列“四大天王”:Rabbit、Rocket、Kafka、Pulsar巅峰对决
本文对比了 RabbitMQ、RocketMQ、Kafka 和 Pulsar 四种消息队列系统,涵盖架构、性能、可用性和适用场景。RabbitMQ 以灵活路由和可靠性著称;RocketMQ 支持高可用和顺序消息;Kafka 专为高吞吐量和低延迟设计;Pulsar 提供多租户支持和高可扩展性。性能方面,吞吐量从高到低依次为
7654 1
|
SQL 存储 大数据
Paimon 在汽车之家的业务实践
本文分享自汽车之家的王刚、范文、李乾⽼师。介绍了汽车之家基于 Paimon 的一些实践,和一些背景。
1192 7
Paimon 在汽车之家的业务实践
|
存储 安全 测试技术
|
机器学习/深度学习 算法 计算机视觉
通过MATLAB分别对比二进制编码遗传优化算法和实数编码遗传优化算法
摘要: 使用MATLAB2022a对比了二进制编码与实数编码的遗传优化算法,关注最优适应度、平均适应度及运算效率。二进制编码适用于离散问题,解表示为二进制串;实数编码适用于连续问题,直接搜索连续空间。两种编码在初始化、适应度评估、选择、交叉和变异步骤类似,但实数编码可能需更复杂策略避免局部最优。选择编码方式取决于问题特性。