赞
踩
MQ 是 Message Queue 的缩写,译为 “消息队列”。
MQ 的主要职责就是转发生产者的消息给消费者,并具备一定的消息缓存能力,在分布式系统中,常常用于各个应用进程之间的通讯行为。
在传统的应用间通讯手段上,往往大多采用直接访问对方URL等同步的数据传输方式,客户端与服务端的消息耦合,这在某些要求实时性和必要性的业务场景下是必需的,但对于某些业务场景,例如短信通知、邮件通知等,本身并不是主业务流程中必要的关键环节,实时性也要求不高,因此,完全可以采用异步的方式来完成,MQ的一个重要作用就是基于这种情况,实现应用间、业务间的异步解耦,是将比较耗时且不需要即时响应的的操作作为消息放入消息队列,同时,由于使用了MQ,只要保证消息格式不变,消息的发送方和接收方并不需要联系彼此,也不需要受对方处理速度的影响,即解耦合。
流量削峰也是MQ 的常用场景,一般在秒杀或团购活动中使用广泛。
目前业界有很多MQ产品,比较出名的有以下这些:
1、ZeroMQ
号称最快的消息队列系统,尤其针对大吞吐量的需求场景。扩展性好,开发比较灵活,采用C语言实现,实际上只是一个socket库的重新封装,如果做为消息队列使用,需要开发大量的代码。ZeroMQ仅提供非持久性的队列,也就是说如果down机,数据将会丢失。
2、RabbitMQ
使用 erlang 语言开发,性能较好,适合企业级开发,但不利于做二次开发和维护。
3、ActiveMQ
历史悠久的 Apache 开源项目。已经在很多产品中得到应用,实现了JMS 1.1 规范,可以和 spring-jms 轻松融合,实现了多种协议,支持持久化,对队列数较多的情况支持不好。
4、RocketMQ
阿里开源的 MQ 组件,由Java 开发,性能很好,使用简单。
5、Kafka
Apache 下的一个子项目,是一个高性能跨语言分布式 Publish/Subscribe 消息队列系统,相对于ActiveMQ 是一个非常轻量级的消息系统,除了性能非常好之外,还是一个工作良好的分布式系统。
Rocket MQ 中重要角色的比喻:
Producer 寄件人、Consumer 收件人、NameServer 邮局、Broker 邮递员。
RocketMQ主要由 Producer、Broker、Consumer 三部分组成,Producer 负责生产消息,Consumer 负责消费消息,Broker 负责存储消息。这三者共同组成 RocketMQ 的消息模型。Broker 在实际部署过程中对应一台服务器,每个 Broker 可以存储多个Topic的消息,每个Topic的消息也可以分片存储于不同的 Broker。Message Queue 用于存储消息的物理地址,每个Topic中的消息地址存储于多个 Message Queue 中。ConsumerGroup 由多个Consumer 实例构成。
生产者负责生产消息,一般由业务系统负责生产消息。一个消息生产者会把业务应用系统里产生的消息发送到broker服务器。RocketMQ提供多种发送方式,同步发送、异步发送、顺序发送、单向发送。同步和异步方式均需要Broker返回确认信息,单向发送不需要。Rocket MQ 支持分布式集群方式部署,Producer 通过MQ的负载均衡模块选择相应的 Broker 集群队列进行消息投递,投递过程支持快速失败并且低延迟。
消费者负责消费消息,一般是后台系统负责异步消费。一个消息消费者会从Broker服务器拉取消息、并将其提供给应用程序。从用户应用的角度而言提供了pull 和 push 两种消费形式。RocketMQ 支持分布式集群方式部署,同时支持实时消息订阅机制和集群、广播等消费模式。
表示一类消息的集合。每个主题包含若干条消息,每条消息只能属于一个主题,是RocketMQ进行消息订阅的基本单位。
Broker Server 消息中转角色,负责存储消息、转发消息。代理服务器在RocketMQ系统中负责接收从生产者发送来的消息并存储、同时为消费者的拉取请求作准备。代理服务器也存储消息相关的元数据,包括消费者组、消费进度偏移和主题和队列消息等。
Name Server 充当路由消息的提供者。生产者或消费者能够通过名字服务查找各主题相应的Broker IP列表。多个Namesrv实例组成集群,但相互独立,没有信息交换。
拉取式:应用通常主动调用Consumer的拉消息方法从Broker服务器拉消息、主动权由应用控制。一旦获取了批量消息,应用就会启动消费过程。
推送式:该模式下Broker收到数据后会主动推送给消费端,该消费模式一般实时性较高。
同一类Producer或Consumer的集合,发送或消费同一类消息且发送或消费逻辑一致。
这是消费者组的两种消息模式。
集群消费模式,相同Consumer Group的每个Consumer实例平均分摊消息。
广播消费模式下,相同Consumer Group的每个Consumer实例都接收全量的消息。
消息系统所传输信息的物理载体,每条消息必须属于一个主题。每个消息拥有唯一的Message ID,且可以携带具有业务标识的Key。系统提供了通过Message ID和Key查询消息的功能。
为消息设置的标志,用于同一主题下区分不同类型的消息。消费者可以根据Tag实现对不同子主题的不同消费逻辑,实现更好的扩展性。
顺序消息分为两种类型。
普通顺序消息,消费者通过同一个消息队列( Topic 分区,称作 Message Queue) 收到的消息是有顺序的,不同消息队列收到的消息则可能是无顺序的。
严格顺序消息,消费者收到的所有消息均是有顺序的。
RocketMQ 架构上主要分为四部分:生产者、消费者、Name Server、Broker Server:
NameServer:NameServer 是一个非常简单的Topic路由注册中心,其角色类似Dubbo中的zookeeper,支持Broker的动态注册与发现。主要包括两个功能:Broker管理,NameServer接受Broker集群的注册信息并且保存下来作为路由信息的基本数据。然后提供心跳检测机制,检查Broker是否还存活;路由信息管理,每个NameServer将保存关于Broker集群的整个路由信息和用于客户端查询的队列信息。然后Producer和Conumser通过NameServer就可以知道整个Broker集群的路由信息,从而进行消息的投递和消费。NameServer通常也是集群的方式部署,各实例间相互不进行信息通讯。Broker是向每一台NameServer注册自己的路由信息,所以每一个NameServer实例上面都保存一份完整的路由信息。当某个NameServer因某种原因下线了,Broker仍然可以向其它NameServer同步其路由信息,Producer,Consumer仍然可以动态感知Broker的路由的信息。
BrokerServer:Broker主要负责消息的存储、投递和查询以及服务高可用保证,为了实现这些功能,Broker包含了以下几个重要子模块。
– Remoting Module:整个Broker的实体,负责处理来自clients端的请求。
– Client Manager:负责管理客户端(Producer/Consumer)和维护Consumer的Topic订阅信息
– Store Service:提供方便简单的API接口处理消息存储到物理硬盘和查询功能。
– HA Service:高可用服务,提供Master Broker 和 Slave Broker之间的数据同步功能。
– Index Service:根据特定的Message key对投递到Broker的消息进行索引服务,以提供消息的快速查询。
1、启动NameServer,NameServer起来后监听端口,等待Broker、Producer、Consumer连上来,相当于一个路由控制中心。
2、Broker启动,跟所有的NameServer保持长连接,定时发送心跳包。心跳包中包含当前Broker信息(IP+端口等)以及存储所有Topic信息。注册成功后,NameServer集群中就有Topic跟Broker的映射关系。
3、收发消息前,先创建Topic,创建Topic时需要指定该Topic要存储在哪些Broker上,也可以在发送消息时自动创建Topic。
4、Producer发送消息,启动时先跟NameServer集群中的其中一台建立长连接,并从NameServer中获取当前发送的Topic存在哪些Broker上,轮询从队列列表中选择一个队列,然后与队列所在的Broker建立长连接从而向Broker发消息。
5、Consumer跟Producer类似,跟其中一台NameServer建立长连接,获取当前订阅Topic存在哪些Broker上,然后直接跟Broker建立连接通道,开始消费消息。
RocketMQ是阿里开源的分布式消息中间件,现在是Apache 的一个顶级项目,在阿里内部使用非常广泛,已经经过了“双十一”万亿级场景下的消息流转。
(待补充)
准备环境:Linux系统 CentOS6 64位(虚拟机),IP:192.168.1.140,JDK8
1、下载并上传 RocketMQ
打开官网下载页面:Release Notes - Apache RocketMQ - Version 4.9.1 下载 Binary 版本 zip 包:
下载后上传到 Linux /usr/local/src/ 目录下:
2、解压缩,并移动到安装目录
unzip rocketmq-all-4.9.1-bin-release.zip
mv rocketmq-all-4.9.1-bin-release /usr/local/rocketmq
3、启动 RocketMQ
切换到 RocketMQ 安装目录,启动 NameServer、BrokerServer,启动脚本在 bin 目录下。& 代表后台启动
nohup ./bin/mqnamesrv &
查看 rocketmq 启动日志,可以看到 The Name Server boot success 字眼,说明NameServer启动成功:
启动 Broker 之前,需要修改几项配置。
# 编辑 bin/runbroker.sh 和 bin/runserver.sh 文件,修改里面的堆大小,视情况而定。
JAVA_OPT="${JAVA_OPT} -server -Xms256m -Xmx256m
JAVA_OPT="${JAVA_OPT} -server -Xms256m -Xmx256m -Xmn128m
然后启动 Broker
# -n 代表 NameServer 地址
nohup ./bin/mqbroker -n localhost:9876 &
查看 Broker 启动日志
4、测试RocketMQ
官方提供了两个测试脚本用于验证 RocketMQ 的可用性。
开启两个终端,分别执行以下命令:
Producer 发送消息:
export NAMESRV_ADDR=localhost:9876
./bin/tools.sh org.apache.rocketmq.example.quickstart.Producer
Consumer 接收消息:
export NAMESRV_ADDR=localhost:9876
./bin/tools.sh org.apache.rocketmq.example.quickstart.Consumer
可以看到发送消息和接收消息都正常完成了:
5、关闭 RocketMQ
bin/mqshutdown broker
bin/mqshutdown namesrv
1、下载
在 git 上下载工程
https://github.com/apache/rocketmq-externals/releases
2、修改配置文件
修改 rocketmq-console\src\main\resources\application.properties
server.port=7777
rocketmq.config.namesrvAddr=192.168.1.140:9876 # nameserver 地址,注意防火墙要开启 9876 端口
3、打成 jar 包,并启动
进入控制台项目,将工程打成 jar 包
mvn clean package -Dmaven.test.skip=true
java -jar target/rocketmq-console-ng-1.0.0.jar
4、访问控制台
打开浏览器,输入 http://localhost:7777 ,就可以看到如下界面:
本节使用 main 方法实现简单的 rocketmq 的消息发送和接收,在此之前,需要确认好是否已经完成前面的 RocketMQ 的环境部署,以及控制台的安装。
1、引入依赖
在需要使用 rocketmq 的项目中加入maven依赖
<dependency>
<groupId>org.apache.rocketmq</groupId>
<artifactId>rocketmq-spring-boot-starter</artifactId>
<version>2.0.2</version>
</dependency>
2、编写消息发送端
import org.apache.rocketmq.client.producer.DefaultMQProducer; import org.apache.rocketmq.client.producer.SendResult; import org.apache.rocketmq.common.message.Message; public class RocketMQSendTest { public static void main(String[] args) throws Exception { // 1. 创建消息生产者,并设置生产组名 DefaultMQProducer producer = new DefaultMQProducer("shop-order"); // 2. 为生产者设置 NameServer 地址 producer.setNamesrvAddr("192.168.1.140:9876"); // 3. 启动生产者 producer.start(); // 4. 构建消息对象,主要是设置消息的主题、标签、内容 Message message = new Message("topic-order", "morty", "Test RocketMQ Message".getBytes()); // 5. 发送消息,第二个参数代表超时时间 SendResult result = producer.send(message, 10000); System.out.println("发送结果:" + result); // 6. 关闭生产者 producer.shutdown(); } }
代码可直接运行,由于 MQ是一种解耦组件,所以,可以直接向MQ 中发送消息而不需要等待消费者。
3、编写消息接收端
消息接收者基于订阅监听机制,需要注册相应的监听器完成消息的消费:
import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer; import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyContext; import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus; import org.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently; import org.apache.rocketmq.common.message.MessageExt; public class RocketMQReceiveTest { public static void main(String[] args) throws Exception { // 1. 创建消息消费者,指定消费者组名 DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("shop-order-consumer"); // 2. 指定 NameServer consumer.setNamesrvAddr("192.168.1.140:9876"); // 3. 指定消费者订阅的主题和标签 consumer.subscribe("topic-order", "*"); // 4. 设置回调函数,并编写消息消费逻辑 consumer.registerMessageListener(new MessageListenerConcurrently() { @Override public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> list, ConsumeConcurrentlyContext consumeConcurrentlyContext) { // 消息消费逻辑 System.out.println("message ====> " + list); // 返回消费成功状态 return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; } }); // 5. 启动消费者 consumer.start(); System.out.println("消费者启动成功!"); } }
启动消费者代码,观察日志:
同时,也可以看到 RocketMQ 控制台有相关主题信息展示:
以 shop-order、shop-user 两个微服务为基础,实现一个下单的消息通知功能。
下单消息通知功能要求,下单后向用户发送下单消息,结构如下图所示:
1、在 shop-order 中添加 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>
2、添加配置
rocketmq:
name-server: 192.168.1.140:9876
producer:
group: shop-order
3、编写发送消息代码
在下单成功后,使用 rocketmq 的接口实现消息的发送。
@Slf4j @RestController @RequestMapping("/order") public class OrderController { @Autowired private RocketMQTemplate rocketMQTemplate; @GetMapping("/prod/{pid}") public Order order(@PathVariable("pid") Integer pid) { log.info("接收到{}号商品的下单请求,准备调用商品微服务", pid); // 调用商品微服务,查询商品信息(略) // 下单(即创建订单并保存)(略) // 订单入库 orderService.createOrder(order); log.info("创建订单成功,订单信息为:{}", JSON.toJSONString(order)); // 使用 rocketMQTemplate 发送下单成功消息 rocketMQTemplate.convertAndSend("order-topic", order); return order; } }
4、测试
访问下单接口,观察 RocketMQ 控制台。
1、添加必要的依赖
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
<dependency>
<groupId>com.alibaba.cloud</groupId>
<artifactId>spring-cloud-alibaba-nacos-discovery</artifactId>
</dependency>
<dependency>
<groupId>org.apache.rocketmq</groupId>
<artifactId>rocketmq-spring-boot-starter</artifactId>
<version>2.0.2</version>
</dependency>
2、配置 rocketmq NameServer 地址
rocketmq:
name-server: 192.168.1.140:9876
3、编写 MQ 监听器
import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
import org.apache.rocketmq.spring.core.RocketMQListener;
@Slf4j
@Service
@RocketMQMessageListener(consumerGroup = "shop-user", topic = "order-topic")
public class OrderListener implements RocketMQListener<Order> {
@Override
public void onMessage(Order order) {
log.info("接收到了下单成功的消息,{}", order);
}
}
4、测试消息接收
启动 shop-user 模块,观察日志,可以看到应用一启动成功,就收到了来自 MQ 的订阅消息:
Copyright © 2003-2013 www.wpsshop.cn 版权所有,并保留所有权利。