RabbitMQ其他知识点
RabbitMQ幂等性
概念
用户对于同一操作发起的一次请求或者多次请求的结果是一致的,不会因为多次点击而产生了副作用。举个最简单的例子,那就是支付,用户购买商品后支付,支付扣款成功,但是返回结果的时候网络异常,此时钱已经扣了,用户再次点击按钮,此时会进行第二次扣款,返回结果成功,用户查询余额发现多扣钱了,流水记录也变成了两条。在以前的单应用系统中,我们只需要把数据操作放入事务中即可,发生错误立即回滚,但是再响应客户端的时候也有可能出现网络中断或者异常等等
重复消费者
消费者在消费MQ中的消息时,MQ已把消息发送给消费者,消费者在给MQ返回ack时网络中断,故MQ未收到确认信息,该条消息会重新发给其他的消费者,或者在网络重连后再次发送给该消费者,但实际上该消费者已成功消费了该条消息,造成消费者消费了重复的消息。
解决思路
MQ消费者的幂等性的解决一般使用全局ID或者写个唯一标识比如时间戳或者UUID或者订单消费者消费MQ中的消息也可利用MQ的该id来判断,或者可按自己的规则生成一个全局唯一id,每次消费消息时用该id先判断该消息是否已消费过。
消费端的幂等性保障
在海量订单生成的业务高峰期,生产端有可能就会重复发生了消息,这时候消费端就要实现幂等性,这就意味着我们的消息永远不会被消费多次,即使我们收到了一样的消息。业界主流的幂等性有两种操作:a.唯一ID+指纹码机制,利用数据库主键去重, b.利用redis.的原子性去实现
唯一ID+指纹码机制
指纹码:我们的一些规则或者时间戳加别的服务给到的唯一信息码;它并不一定是我们系统生成的,基本都是由我们的业务规则拼接而来,但是一定要保证唯一性,然后就利用查询语句进行判断这个id是否存在数据库中,优势就是实现简单就一个拼接,然后查询判断是否重复;劣势就是在高并发时,如果是单个数据库就会有写入性能瓶颈当然也可以采用分库分表提升性能,但也不是我们最推荐的方式。
Redis原子性
利用redis.,执行setnx,命令,天然具有幂等性。从而实现不重复消费
RabbitMQ优先级队列
使用场景
在我们系统中有一个订单催付的场景,我们的客户在天猫下的订单,淘宝会及时将订单推送给我们,如果在用户设定的时间内未付款那么就会给用户推送一条短信提醒,很简单的一个功能对吧,但是,tmall商家对我们来说,肯定是要分大客户和小客户的对吧,比如像苹果,小米这样大商家一年起码能给我们创造很大的利润,所以理应当然,他们的订单必须得到优先处理,而曾经我们的后端系统是使用redis.来存放的定时轮询,大家都知道redis.只能用List做一个简简单单的消息队列,并不能实现一个优先级的场景,所以订单量大了后采用RabbitMQ进行改造和优化,如果发现是大客户的订单给一个相对比较高的优先级,否则就是默认优先级。

package com.yuan.rabbitmq.simple;
import com.rabbitmq.client.AMQP;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import java.io.IOException;
import java.util.HashMap;
import java.util.concurrent.TimeoutException;
/**
* 简单队列生产者
*
* @author 111
*/
public class Producer {
public static void main(String[] args) {
//1.创建连接工厂
ConnectionFactory connectionFactory = new ConnectionFactory();
//2.设置连接属性
connectionFactory.setHost("120.25.213.140");
connectionFactory.setPort(5672);
connectionFactory.setVirtualHost("/");
connectionFactory.setUsername("admin");
connectionFactory.setPassword("admin");
Connection connection = null;
Channel channel = null;
try {
//3.从连接工厂获取
connection = connectionFactory.newConnection("生产者");
//4.从连接中获取通道channel
channel = connection.createChannel();
//5.声明队列queue存贮消息
String queueName = "queue1";
HashMap<String, Object> arguments = new HashMap<>();
//官方允许是0-255之间﹑此处设置10允许优化级范围为0-10不要设置过大, 浪费CPU内存
arguments.put("x-max-priority",10);
channel.queueDeclare(queueName, true, false, false, arguments);
//6.准备发送的消息内容
for (int i = 1; i < 11; i++) {
String msg = "info"+i;
//如果i等于5设置优先级为5
if (i==5) {
AMQP.BasicProperties properties=
new AMQP.BasicProperties().builder().priority(5).build();
channel.basicPublish("", queueName, properties, msg.getBytes());
}else {
channel.basicPublish("", queueName, null, msg.getBytes());
}
}
System.out.println("消息发送成功");
} catch (IOException | TimeoutException e) {
e.printStackTrace();
} finally {
//7.释放连接关闭通道
if (channel != null && channel.isOpen()) {
try {
channel.close();
} catch (IOException | TimeoutException e) {
e.printStackTrace();
}
}
if (connection != null) {
try {
connection.close();
} catch (IOException e) {
e.printStackTrace();
}
}
}
}
}
package com.yuan.rabbitmq.simple;
import com.rabbitmq.client.*;
import java.io.IOException;
import java.util.concurrent.TimeoutException;
/**
* 服务消费者
*/
public class Consumer {
public static void main(String[] args) {
//1.创建连接工厂
ConnectionFactory connectionFactory = new ConnectionFactory();
//2.设置连接属性
connectionFactory.setHost("120.25.213.140");
connectionFactory.setPort(5672);
connectionFactory.setVirtualHost("/");
connectionFactory.setUsername("admin");
connectionFactory.setPassword("admin");
Connection connection = null;
Channel channel = null;
try {
//3.从连接工厂获取
connection = connectionFactory.newConnection("生产者");
//4.从连接中获取通道channel
channel = connection.createChannel();
//服务消费者,消费
channel.basicConsume("queue1", true, new DeliverCallback() {
@Override
public void handle(String consumerTag, Delivery message) throws IOException {
System.out.println("收到的消息是" + new String(message.getBody(), "utf-8"));
}
//如果又异常接收失败
}, new CancelCallback() {
@Override
public void handle(String consumerTag) throws IOException {
System.out.println("接收失败");
}
});
System.out.println("开始接收消息");
//不关闭阻塞
System.in.read();
} catch (IOException | TimeoutException e) {
e.printStackTrace();
} finally {
//7.释放连接关闭通道
if (channel != null && channel.isOpen()) {
try {
channel.close();
} catch (IOException | TimeoutException e) {
e.printStackTrace();
}
}
if (connection != null) {
try {
connection.close();
} catch (IOException e) {
e.printStackTrace();
}
}
}
}
}
RabbitMQ惰性队列
使用场景
RabbitMQ从3.6.0版本开始引入了惰性队列的概念。惰性队列会尽可能的将消息存入磁盘中,而在消费者消费到相应的消息时才会被加载到内存中,它的一个重要的设计目标是能够支持更长的队列,即支持更多的消息存储。当消费者由于各种各样的原因(比如消费者下线、宕机亦或者是由于维护而关闭等)而致使长时间内不能消费消息造成堆积时,惰性队列就很有必要了。
默认情况下,当生产者将消息发送到RabbitMQ的时候,队列中的消息会尽可能的存储在内存之中,这样可以更加快速的将消息发送给消费者。即使是持久化的消息,在被写入磁盘的同时也会在内存中驻留一份备份。当RabbitMQ需要释放内存的时候,会将内存中的消息换页至磁盘中,这个操作会耗费较长的时间,也会阻塞队列的操作,进而无法接收新的消息。虽然RabbitMQ.的开发者们一直在升级相关的算法,但是效果始终不太理想,尤其是在消息量特别大的时候。
两种模式
队列具备两种模式: default和lazy。默认的为default模式,在3.6.0之前的版本无需做任何变更。lazy模式即为惰性队列的模式,可以通过调用channel.queueLecare 云日RVix在今女重言的优先级Policy的方式设置,如果一个队列同时使用这两种方式设置的话,那么Policy的方式具备更高的优先级。如果要通过声明的方式改变已有队列的模式的话,那么只能先删除队列,然后再重新声明一个新的在队列声明的时候可以通过"x-queue-mode"参数来设置队列的模式,取值为"default"和"lazy"。下面示例中演示了一个惰性队列的声明细节:
Java程序设置
Map<String, Object> args = new HashMap<String, Object>();
args.put("x-queue-mode" , "lazy");
channel.queueDeclare("myqueue", false, false, false, args);
网页模式设置

内存开销对比

在发送1百万条消息,每条消息大概占1KB的情况下,普通队列占用内存是1.2GB,而惰性队列仅仅占用1.5MB
RabbitMQ集群(单机多实例搭建)
RabbitMQ这款消息队列中间件产品本身是基于Erlang编写,Erlang语言天生具备分布式特性(通过同步Erlang集群各节点的magic cookie来实现)。因此,RabbitMQ天然支持Clustering。这使得RabbitMQ本身不需要像ActiveMQ、Kafka那样通过ZooKeeper分别来实现HA方案和保存集群的元数据。集群是保证可靠性的一种方式,同时可以通过水平扩展以达到增加消息吞吐量能力的目的。 在实际使用过程中多采取多机多实例部署方式,为了便于同学们练习搭建,有时候你不得不在一台机器上去搭建一个rabbitmq集群,本章主要针对单机多实例这种方式来进行开展。
主要参考官方文档:https://www.rabbitmq.com/clustering.html
环境准备

安装
#安装erlang
[root@yuan rabbitmq]# rpm -Uvh erlang-solutions-2.0-1.noarch.rpm
[root@yuan rabbitmq]# yum install -y erlang
#查看是否安装成功
[root@yuan rabbitmq]# erl -v
#安装socat
[root@yuan rabbitmq]# yum install -y socat
#安装rabbitmq
[root@yuan rabbitmq]#rpm -Uvh rabbitmq-server-3.9.13-1.el7.noarch.rpm
#如果报错误:依赖检测失败:使用下面这段
[root@yuan rabbitmq]# rpm -Uvh rabbitmq-server-3.9.13-1.el7.noarch.rpm --force --nodeps
[root@yuan rabbitmq]# yum install rabbitmq-server -y
已安装上面的可以忽略
配置的前提是你的rabbitmq可以运行起来,比如”ps aux|grep rabbitmq”你能看到相关进程,又比如运行“rabbitmqctl status”你可以看到类似如下信息,而不报错:
ps aux|grep rabbitmq

或者
systemctl status rabbitmq-server
注意:确保RabbitMQ可以运行的,确保完成之后,把单机版的RabbitMQ服务停止,后台看不到RabbitMQ的进程为止
单机多实例搭建
**场景:**假设有两个rabbitmq节点,分别为rabbit-1, rabbit-2,rabbit-1作为主节点,rabbit-2作为从节点。 启动命令:RABBITMQ_NODE_PORT=5672 RABBITMQ_NODENAME=rabbit-1 rabbitmq-server -detached 结束命令:rabbitmqctl -n rabbit-1 stop
第一步:启动第一个节点rabbit-1
[root@yuan rabbitmq]# sudo RABBITMQ_NODE_PORT=5672 RABBITMQ_NODENAME=rabbit-1 rabbitmq-server start &
启动第二个节点rabbit-2
[root@yuan rabbitmq]# sudo RABBITMQ_NODE_PORT=5673 RABBITMQ_SERVER_START_ARGS="-rabbitmq_management listener [{port,15673}]" RABBITMQ_NODENAME=rabbit-2 rabbitmq-server start &
启动第三个节点rabbit-3
[root@yuan rabbitmq]# sudo RABBITMQ_NODE_PORT=5674 RABBITMQ_SERVER_START_ARGS="-rabbitmq_management listener [{port,15674}]" RABBITMQ_NODENAME=rabbit-3 rabbitmq-server start &
查看是否启动成功:ps aux|grep rabbitmq

rabbit-1操作作为主节点
#停止应用
sudo rabbitmqctl -n rabbit-1 stop_app
#目的是清除节点上的历史数据(如果不清除,无法将节点加入到集群)
sudo rabbitmqctl -n rabbit-1 reset
#启动应用
sudo rabbitmqctl -n rabbit-1 start_app
rabbit2操作为从节点
# 停止应用
sudo rabbitmqctl -n rabbit-2 stop_app
# 目的是清除节点上的历史数据(如果不清除,无法将节点加入到集群)
sudo rabbitmqctl -n rabbit-2 reset
# 将rabbit2节点加入到rabbit1(主节点)集群当中【yuan服务器的主机名】
sudo rabbitmqctl -n rabbit-2 join_cluster rabbit-1@yuan
# 启动应用
sudo rabbitmqctl -n rabbit-2 start_app
rabbit3操作为从节点
# 停止应用
sudo rabbitmqctl -n rabbit-3 stop_app
# 目的是清除节点上的历史数据(如果不清除,无法将节点加入到集群)
sudo rabbitmqctl -n rabbit-3 reset
# 将rabbit2节点加入到rabbit1(主节点)集群当中【yuan服务器的主机名】
sudo rabbitmqctl -n rabbit-3 join_cluster rabbit-1@yuan
# 启动应用
sudo rabbitmqctl -n rabbit-3 start_app
验证集群状态
sudo rabbitmqctl cluster_status -n rabbit-1

RabbitMQ集群Web界面
#下载web界面插件
[root@yuan rabbitmq]# rabbitmq-plugins enable rabbitmq_management
注意在访问的时候:web结面的管理需要给15672 rabbit-1 和15673的rabbit-2,15674的rabbit-3 设置用户名和密码。如下:
rabbitmqctl -n rabbit-1 add_user admin admin
rabbitmqctl -n rabbit-1 set_user_tags admin administrator
rabbitmqctl -n rabbit-1 set_permissions -p / admin ".*" ".*" ".*"
rabbitmqctl -n rabbit-2 add_user admin admin
rabbitmqctl -n rabbit-2 set_user_tags admin administrator
rabbitmqctl -n rabbit-2 set_permissions -p / admin ".*" ".*" ".*"
rabbitmqctl -n rabbit-3 add_user admin admin
rabbitmqctl -n rabbit-3 set_user_tags admin administrator
rabbitmqctl -n rabbit-3 set_permissions -p / admin ".*" ".*" ".*"

小结
如果采用多机部署方式,需读取其中一个节点的cookie, 并复制到其他节点(节点之间通过cookie确定相互是否可通信)。cookie存放在/var/lib/rabbitmq/.erlang.cookie。 例如:主机名分别为rabbit-1、rabbit-2 1、逐个启动各节点 2、配置各节点的hosts文件( vim /etc/hosts) ip1:rabbit-1 ip2:rabbit-2 其它步骤雷同单机部署方式
RabbitMQ集群(多机多实例搭建)
搭建步骤
-
修改3台机器的主机名称
vim /etc/hostname
⒉配置各个节点的hosts文件,让各个节点都能互相识别对方,三台机器都复制一份:对应的三台机器的Ip地址
vim/etc/hosts
-
211.55.74 node1
-
211.55.75 node2
-
211.55.76 node3
-
以确保各个节点的cookie文件使用的是同一个值
在node1上执行远程操作命令
scp /var/lib/rabbitmq/.erlang.cookie root@node2:/var/lib/rabbitmq/.erlang.cookie scp /var/lib/rabbitmq/.erlang.cookie root@node3:/var/lib/rabbitmq/.erlang.cookie
- 启动RabbitMQ服务,顺带启动Erlang虚拟机和RbbitMQ应用服务(在三台节点上分别执行以下命令)
rabbitmq-server -detached
-
在节点2执行
rabbitmqctl stop_app
(rabbitmqctl stop 会将Erlang 虚拟机关闭,rabbitmqctl stop_app 只关闭 RabbitMQ 服务)
rabbitmqctl reset
rabbitmqctl join_cluster rabbit@node1
rabbitmqctl start_app(只启动应用服务)
- 在节点3执行
rabbitmqctl stop_app rabbitmqctl reset rabbitmqctl join_cluster rabbit@node2 rabbitmqctl start_app
-
集群状态
rabbitmqctl cluster_status
-
需要重新设置用户
创建账号
rabbitmqctl add_user admin 123
设置用户角色
rabbitmqctl set_user_tags admin administrator
设置用户权限
rabbitmqctl set_permissions -p "/" admin "." "." ".*"
-
解除集群节点(解除集群使用 node2和node3机器分别执行)
rabbitmqctl stop_app
abbitmqctl reset
rabbitmqctl start_app
rabbitmqctl cluster_status
rabbitmqctl forget_cluster_node rabbit@node2(node1 机器上执行)
RabbitMQ镜像队列
注:多机多实例搭建才有用
使用镜像的原因
如果RabbitMQ集群中只有一个Broker 节点,那么该节点的失效将导致整体服务的临时性不可用,并且也可能会导致消息的丢失。可以将所有消息都设置为持久化,并且对应队列的durable属性也设置为true,但是这样仍然无法避免由于缓存导致的问题:因为消息在发送之后和被写入磁盘井执行副盘动作之间存在一个短暂却会产生问题的时间窗。通过publisherconfirm机制能够确保客户端知道哪些消息己经存入磁盘,尽管如此,一般不希望遇到因单点故障导致的服务不可用。 引入镜像队列(Mirror Queue)的机制,可以将队列镜像到集群中的其他Broker节点之上,如果集群中的一个节点失效了,队列能自动地切换到镜像中的另一个节点上以保证服务的可用性。
1、启动三台集群节点
2、随便找一个节点添加 policy

3、 在 node1 上创建一个队列发送一条消息,队列存在镜像队列

就算整个集群只剩下一台机器了依然能肖费队列里面的消息
说明队列里面的消息被镜像队列传递到相应机器里面了
Haproxy+Keepalive 实现高可用(负载均衡)
注:多机多实例搭建才有用
整体架构图

Haproxy 实现负载均衡
HAProxy 提供高可用性、负载均衡及基于TCPHTTP 应用的代理,支持虚拟主机,它是免费、快速并且可靠的一种解决方案,包括 Twitter,Reddit,StackOverflow,GitHub 在内的多家知名互联网公司在使用。HAProxy 实现了一种事件驱动、单一进程模型,此模型支持非常大的井发连接数。
扩展 nginx,lvs,haproxy 之间的区别: http://www.ha97.com/5646.html
3、搭建步骤 1、下载 haproxy(在 node1 和 node2)
yum -y install haproxy
2、修改 node1 和 node2 的 haproxy.cfg
vim /etc/haproxy/haproxy.cfg
需要修改红色 IP 为当前机器 IP

3、在两台节点启动 haproxy
haproxy -f /etc/haproxy/haproxy.cfg
ps -ef | grep haproxy
4、访问地址
http://10.211.55.71:8888/stats
Keepalived 实现双机(主备)热备
试想如果前面配置的 HAProxy 主机突然宕机或者网卡失效,那么虽然 RbbitMQ 集群没有任何故障但是对于外界的客户端来说所有的连接都会被断开结果将是灾难性的为了确保负载均衡服务的可靠性同样显得十分重要,这里就要引入 Keepalived 它能够通过自身健康检查、资源接管功能做高可用(双机热备),实现故障转移.
1、 下载 keepalived
yum -y install keepalived
2、节点 node1 配置文件
vim /etc/keepalived/keepalived.conf
把资料里面的 keepalived.conf 修改之后替换
3、节点 node2 配置文件
需要修改global_defs 的 router_id,如:nodeB
其次要修改 vrrp_instance_VI 中 state 为"BACKUP"; 最后要将priority 设置为小于 100 的值
4、添加 haproxy_chk.sh
(为了防止 HAProxy 服务挂掉之后 Keepalived 还在正常工作而没有切换到 Backup 上,所以
这里需要编写一个脚本来检测 HAProxy 务的状态,当 HAProxy 服务挂掉之后该脚本会自动重启
HAProxy 的服务,如果不成功则关闭 Keepalived 服务,这样便可以切换到 Backup 继续工作) vim /etc/keepalived/haproxy_chk.sh(可以直接上传文件)
修改权限 chmod 777 /etc/keepalived/haproxy_chk.sh
5、启动 keepalive 命令(node1 和 node2 启动)
systemctl start keepalived
6、观察 Keepalived 的日志
tail -f /var/log/messages -n 200
7、观察最新添加的 vip ip add show
8、 node1 模拟 keepalived 关闭状态
systemctl stop keepalived
9、使用 vip 地址来访问 rabbitmq 集群
Federation Exchange(交换机)
注:多机多实例搭建才有用
使用它的原因
(broker北京),(broker深圳)彼此之间相距甚远,网络延迟是一个不得不面对的问题。有一个在北京的业务(Client北京)需要连接(broker北京),向其中的交换器exchangeA发送消息,此时的网络延迟很小,(Client北京)可以迅速将消息发送至exchangeA.中,就算在开启了publisherconfirm.机制或者事务机制的情况下,也可以迅速收到确认信息。此时又有个在深圳的业务(Client深圳)需要向exchangeA.发送消息,那么(Client深圳)(broker北京)之间有很大的网络延迟,(Client深圳)将发送消息至exchangeA.会经历一定的延迟,尤其是在开启了publisherconfirm.机制或者事务机制的情况下,(Client深圳)会等待很长的延迟时间来接收(broker北京)的确认信息,进而必然造成这条发送线程的性能降低,甚至造成一定程度上的阻塞。
将业务(Client深圳)部署到北京的机房可以解决这个问题,但是如果(Client深圳)调用的另些服务都部署在深圳,那么又会引发新的时延问题,总不见得将所有业务全部部署在一个机房,那么容灾又何以实现?这里使用Federation插件就可以很好地解决这个问题.
搭建步骤
-
需要保证每台节点单独运行
-
在每台机器上开启federation相关插件
rabbitmq-plugins enable rabbitmq_federation
rabbitmq-plugins enable rabbitmq_federation_management

3、原理图(先运行 consumer 在 node2 创建 fed_exchange)

4、在 downstream(node2)配置 upstream(node1)

5、添加policy

Federation Exchange(实现)
package com.yuan.rabbitmq.federation;
import com.rabbitmq.client.*;
import java.io.IOException;
import java.util.concurrent.TimeoutException;
public class Consumer {
//队列的名称
public static final String QUEUE_NAME = "mirrior_hello";
//交换机的名称
public static final String FED_EXCHANGE = "fed_exchange";
public static void main(String[] args) throws IOException, TimeoutException {
//创建链接工厂
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("192.168.200.143");
factory.setUsername("admin");
factory.setPassword("admin");
Connection connection = factory.newConnection();
Channel channel = connection.createChannel();
channel.exchangeDeclare(FED_EXCHANGE, BuiltinExchangeType.DIRECT);
channel.queueDeclare("node2_queue",true,false,false,null);
channel.queueBind("node2_queue",FED_EXCHANGE,"routeKey");
//消息的接收
DeliverCallback deliverCallback = ((consumerTag, message) -> {
System.out.println("接收到的消息" + new String(message.getBody()));
});
//消息取消接收收,执行下面的内容
CancelCallback cancelCallback = (consumerTag -> {
System.out.println(consumerTag + "消息者取消消费接口回调逻辑");
});
/**
* 水消费者消费消息
* 1.消费哪个队列
* 2.消费成功之后是否要自动应答true 代表的自动应答 false 代表手动应答
* 3.消费者未成功消费的回调
* 4.消费者取录消费的回调
*/
System.out.println("等待接收消息.....");
channel.basicConsume(QUEUE_NAME, true, deliverCallback, cancelCallback);
}
}
Federation Queue(队列)
使用它的原因
联邦队列可以在多个Broker节点(或者集群)之间为单个队列提供均衡负载的功能。一个联邦队列可以连接一个或者多个上游队列(upstream queue),并从这些上游队列中获取消息以满足本地消费者消费消息的需求。
搭建步骤
.png)

Shovel
注:多机多实例搭建才有用
使用它的原因
Federation具备的数据转发功能类似,Shovel够可靠、持续地从一个Broker中的队列(作为源端,即source)拉取数据并转发至另一个Broker中的交换器(作为目的端,即destination)。作为源端的队列和作为目的端的交换器可以同时位于同一个Broker,也可以位于不同的Broker上。Shovel可以翻译为"铲子",是一种比较形象的比喻,这个"铲子"可以将消息从一方"铲子"另一方。Shovel行为就像优秀的客户端应用程序能够负责连接源和目的地、负责消息的读写及负责连接失败问题的处理。
搭建步骤
启用shovel插件命令:
rabbitmq-plugins enable rabbitmq_shovel rabbitmq-plugins enable rabbitmq_shovel_management

⒉.原理图(在源头发送的消息直接回进入到目的地队列)

- 使用shovel
