一、RabbitMQ 概述
1 消息队列
消息(Message)是指在应用间传送的数据。消息可以非常简单,比如只包含文本字符串,也可以更复杂,可能包含嵌入对象。不建议传递对象,如果需要传递复杂数据建议传递Json。
消息队列(Message Queue)是一种应用间的通信方式,消息发送后可以立即返回,由消息系统来确保消息的可靠传递。消息发布者只管把消息发布到 MQ 中而不用管谁来取,消息使用者只管从 MQ 中取消息而不管是谁发布的。这样发布者和使用者都不用知道对方的存在。
消息队列用于业务解耦、最终一致性、广播、错峰流控等等情况
2 RabbitMQ 特点
RabbitMQ 是一个由 Erlang 语言开发的 AMQP 的开源实现。
AMQP :Advanced Message Queue,高级消息队列协议。它是应用层协议的一个开放标准,为面向消息的中间件设计,基于此协议的客户端与消息中间件可传递消息,并不受产品、开发语言等条件的限制。
RabbitMQ 最初起源于金融系统,用于在分布式系统中存储转发消息,在易用性、扩展性、高可用性等方面表现不俗。具体特点包括:
二 RabbitMQ的消息发送和接收机制
1 概述
所有 MQ 产品从模型抽象上来说都是一样的过程: 消费者(consumer)订阅某个队列。生产者(producer)创建消息,然后发布到队列(queue)中, 最后将消息发送到监听的消费者。
RabbitMQ的内部接收如下:
2 AMQP 中的消息路由
AMQP 中消息的路由过程和 Java 开发者熟悉的 JMS 存在一些差别,AMQP 中增加了 Exchange 和 Binding 的角色。生产者把消息发布到 Exchange 上,消息最终到达队列并被消费者接收,而 Binding 决定交换器的消息应该发送到那个队列
3 Exchange 类型
Exchange分发消息时根据类型的不同分发策略有区别,目前共四种类型:direct、fanout、topic、headers 。
三、Java RabbitMQ
1 基础操作
依赖
<dependencies>
<dependency>
<groupId>com.rabbitmq</groupId>
<artifactId>amqp-client</artifactId>
<version>5.1.1</version>
</dependency>
</dependencies>
编写消息发送类
// 创建连接
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("192.168.174.135");
factory.setPort(5672);
factory.setUsername("root");
factory.setPassword("root");
factory.setVirtualHost("/");
try (Connection connection = factory.newConnection(); Channel channel = connection.createChannel()) {
/*
* 定义队列
* 参数1:队列名
* 参数2:是否持久化
* 参数3:是否排他(当有一个消费者监听时,是否还可以让其他消费者监听)
* 参数4:是否自动删除(如果没有任何消费者监听这个队列,是否要删除队列)
* 参数5:属性,填null即可
* */
channel.queueDeclare("myQueue", true, false, false, null);
/*
* 定义交换机
* 参数1:交换机名字
* 参数2:交换机类型
* 参数3:是否持久化
* */
channel.exchangeDeclare("myExchange", "direct", true);
/*
* 绑定队列
* 参数1:队列名字
* 参数2:交换机名字
* 参数3:routing-key(路由键)
* */
channel.queueBind("myQueue", "myExchange", "myKey");
String message = "hello mq";
/*
* 发送消息
* 参数1:交换机名称
* 参数2:路由键
* 参数3:属性,null即可
* 参数4:消息内容
* */
channel.basicPublish("myExchange", "myKey", null, message.getBytes(StandardCharsets.UTF_8));
} catch (IOException | TimeoutException e) {
e.printStackTrace();
}
编写消息接收类
// 创建连接
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("192.168.174.135");
factory.setPort(5672);
factory.setUsername("root");
factory.setPassword("root");
factory.setVirtualHost("/");
Connection connection = null;
Channel channel = null;
try {
connection = factory.newConnection();
channel = connection.createChannel();
channel.queueDeclare("myQueue", true, false, false, null);
channel.exchangeDeclare("myExchange", "direct", true);
channel.queueBind("myQueue", "myExchange", "myKey");
/*
* 监听接收消息
* 参数1:队列名
* 参数2:是否自动确认
* 参数3:回调函数
*/
channel.basicConsume("myQueue", true, new DefaultConsumer(channel) {
/**
* 接收消息的回调函数
* @param consumerTag 消费者编号
* @param envelope 消息的基础属性
* @param properties 基础消息的属性
* @param body 消息内容
*/
@Override
public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
System.out.println(new String(body, StandardCharsets.UTF_8));
}
});
} catch (IOException | TimeoutException e) {
e.printStackTrace();
}
2 事务消息
事务消息与数据库的事务类似,只是MQ中的消息是要保证消息是否会全部发送成功,防止丢失消息的一种策略。
RabbitMQ有两种方式来解决这个问题:
事务的实现主要是对信道(Channel)的设置,主要的方法有三个:
3 发送者确认模式
Confirm发送方确认模式使用和事务类似,也是通过设置Channel进行发送方确认的,最终达到确保所有的消息全部发送成功
channel.confirmsSelect(); // 开启发送者确认模式
channel.waitForConfirms(5000L); // 确认是否发送成功
waitForConfirms 方法会判定在一定时间内,是否成功发送消息,如果成功返回 true,false则失败。如果抛出中断异常,那么不确定有没有发送成功,需要补发信息。
channel.addConfirmListener(new ConfirmListener() {
//消息确认收到后回调的方法
public void handleAck(long l, boolean b) throws IOException {
System.out.println("收到消息 编号:" + l + " 是否批量:" + b);
}
//消息确认没有收到后的回调方法
public void handleNack(long l, boolean b) throws IOException {
System.out.println("没有收到消息 编号:" + l + " 是否批量:" + b);
}
});
addConfirmListener 是异步确认,他的参数,需要定义接收成功和失败的回调函数。
4 消费者确认模式
消费者在声明队列时,可以指定 noAck 参数,当 noAck=false 时,RabbitMQ会等待消费者显式发回 ack 信号后才从内存(和磁盘,如果是持久化消息的话)中移去消息。否则,RabbitMQ会在队列中消息被消费后立即删除它。
在Consumer中Confirm模式中分为手动确认和自动确认。 手动确认主要并使用以下方法:
- basicAck(): 用于肯定确认,multiple参数用于多个消息确认。
- basicRecover():是路由不成功的消息可以使用recovery重新发送到队列中。
- basicReject():是接收端告诉服务器这个消息我拒绝接收,不处理,可以设置是否放回到队列中还是丢掉, 而且只能一次拒绝一个消息,官网中有明确说明不能批量拒绝消息,为解决批量拒绝消息才有了 basicNack。
- basicNack():可以一次拒绝N条消息,客户端可以设置basicNack方法的multiple参数为true。
channel.basicConsume(queueName,false,new DefaultConsumer(channel){
public void handleDelivery(String consumerTag, Envelope
envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
Channel c=this.getChannel();
try{
System.out.println("—–准备处理消息—–");
String message = new String(body);
System.out.println("Receive–" + message);
//获取消息的编号
long msgTag=envelope.getDeliveryTag();
//手动确认消息,需要在所有的操作全部完成后将消息从队列中移除,
//参数 1 为取消确认的消息编号
//参数 2 为是否批量确认true表示批量确认消息,会自动移除小于等于当 前消息编号的所有消息
c.basicAck(msgTag,true);
}catch ( Exception e){
//将消息重新放回队列,如果消息处理出现了异常则将消息从新放回队列中,尝试再次处理消息
c.basicRecover();
}
}
});
四、SpringBoot集成RabbitMQ
1 配置
依赖
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
</dependency>
配置文件
spring:
rabbitmq:
host: localhost
port: 5672
username: root
password: root
配置类
@Configuration
public class RabbitCollectConfig {
@Bean
public Queue queue() {
return new Queue("bootQueue", true, false, false, null);
}
@Bean
public DirectExchange directExchange() {
return new DirectExchange("bootExchange", true, false);
}
@Bean
public Binding binding(Queue queue, Exchange exchange) {
return new Binding("bootQueue", Binding.DestinationType.QUEUE, exchange.getName(), "bootExchange", null);
}
}
2 生产者
@Autowired
private AmqpTemplate amqpTemplate;
@Test
void send() {
amqpTemplate.convertAndSend("bootExchange", "bootKey", "test");
}
3 消费者
直接获取
@Autowired
private AmqpTemplate amqpTemplate;
@Test
void receive() {
amqpTemplate.receiveAndConvert("bootQueue");
}
或者监听
@Service
public class MessageService {
@RabbitListener
public void receiveMessage(String message) {
System.out.println(message);
}
}
五、使用 Canal 框架同步数据
添加依赖
<dependency>
<groupId>top.javatool</groupId>
<artifactId>canal-spring-boot-starter</artifactId>
<version>1.2.1-RELEASE</version>
</dependency>
配置文件
canal:
server: Canal服务部署的地址:11111
destination: example
user-name: canal
password: Canal_2020
logging:
level:
root: info
top:
javatool:
canal:
client:
client:
AbstractCanalClient: error
添加 handler
@Slf4j
@Component
@CanalTable(value = "t_order_info")
public class OrderaInfoHandler implements EntryHandler<OrderInfo> {
@Autowired
private StringRedisTemplate redisTemplate;
@Override
public void insert( OrderInfo orderInfo) {
log.info("当有数据插入的时候会触发这个方法");
}
@Override
public void update(OrderInfo before, OrderInfo after) {
log.info("当有数据更新的时候会触发这个方法");
}
@Override
public void delete(OrderInfo orderInfo) {
log.info("当有数据删除的时候会触发这个方法");
}
}
编写实体类 OrderInfo,当指定的表修改之后,即可触发方法,可以发送MQ、缓存、同步其他中间件等。





