上一篇消息中间件——RabbitMQ(七)高级特性全在这里!(上)中咱们介绍了消息如何保障100%的投递成功?
,幂等性概念详解
,在海量订单产生的业务高峰期,如何避免消息的重复消费的问题?
,Confirm确认消息、Return返回消息
。这篇咱们来介绍下下面内容。java
咱们通常就在代码中编写while循环,进行consumer.nextDelivery方法进行获取下一条消息,而后进行消费处理!git
可是这种轮训的方式确定是很差的,代码也比较low。github
/** * * @ClassName: Producer * @Description: 生产者 * @author Coder编程 * @date2019年7月30日 下午23:15:51 * */
public class Producer {
public static void main(String[] args) throws Exception {
//1 建立ConnectionFactory
Connection connection = ConnectionUtils.getConnection();
Channel channel = connection.createChannel();
String exchange = "test_consumer_exchange";
String routingKey = "consumer.save";
String msg = "Hello RabbitMQ Consumer Message";
for(int i =0; i<5; i ++){
channel.basicPublish(exchange, routingKey, true, null, msg.getBytes());
}
}
}
复制代码
/** * * @ClassName: Consumer * @Description: 消费者 * @author Coder编程 * @date2019年7月30日 下午23:13:51 * */
public class Consumer {
public static void main(String[] args) throws Exception {
// 建立ConnectionFactory
Connection connection = ConnectionUtils.getConnection();
Channel channel = connection.createChannel();
String exchangeName = "test_consumer_exchange";
String routingKey = "consumer.#";
String queueName = "test_consumer_queue";
channel.exchangeDeclare(exchangeName, "topic", true, false, null);
channel.queueDeclare(queueName, true, false, false, null);
channel.queueBind(queueName, exchangeName, routingKey);
//实现本身的MyConsumer()
channel.basicConsume(queueName, true, new MyConsumer(channel));
}
}
复制代码
/** * * @ClassName: MyConsumer * @Description: TODO * @author Coder编程 * @date 2019年7月30日 下午23:11:55 * */
public class MyConsumer extends DefaultConsumer {
public MyConsumer(Channel channel) {
super(channel);
}
//根据需求,重写本身须要的方法。
@Override
public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
System.err.println("-----------consume message----------");
//消费标签
System.err.println("consumerTag: " + consumerTag);
//这个对象包含许多关键信息
System.err.println("envelope: " + envelope);
System.err.println("properties: " + properties);
System.err.println("body: " + new String(body));
}
}
复制代码
为何不在生产端进行限流呢?面试
由于在高并发的状况下,客户量就是很是大,因此很难在生产端作限制。所以咱们能够用MQ在消费端作限流。编程
参数解释: prefetchSize:0 prefetchCount:会告诉RabbitMQ不要同时给一个消费者推送多于N个消息,即一旦有N个消息尚未ack,则该consumer将block掉,直到有消息ack。 global: true\false 是否将上面设置应用于channel,简单点说,就是上面限制是channel级别仍是consumer级别。 prefetchSize和global这两项,rabbitmq没有实现,暂且不研究prefetch_count在no_ask = false的状况下生效,即在自动应答的状况下这两个值是不生效的。bash
第一个参数:消息的限制大小,消息多少兆。通常不作限制,设置为0 第二个参数:一次最多处理多少条,实际工做中设置为1就好 第三个参数:限流策略在什么上应用。在RabbitMQ通常有两个应用级别:1.通道 2.Consumer级别。通常设置为false,true 表示channel级别,false表示在consumer级别服务器
/** * * @ClassName: Producer * @Description: 生产者 * @author Coder编程 * @date2019年7月30日 下午23:15:51 * */
public class Producer {
public static void main(String[] args) throws Exception {
//1 建立ConnectionFactory
Connection connection = ConnectionUtils.getConnection();
Channel channel = connection.createChannel();
String exchange = "test_qos_exchange";
String routingKey = "qos.save";
String msg = "Hello RabbitMQ QOS Message";
for(int i =0; i<5; i ++){
channel.basicPublish(exchange, routingKey, true, null, msg.getBytes());
}
}
}
复制代码
/** * * @ClassName: Consumer * @Description: 消费者 * @author Coder编程 * @date2019年7月30日 下午23:13:51 * */
public class Consumer {
public static void main(String[] args) throws Exception {
//1 建立ConnectionFactory
Connection connection = ConnectionUtils.getConnection();
Channel channel = connection.createChannel();
String exchangeName = "test_qos_exchange";
String queueName = "test_qos_queue";
String routingKey = "qos.#";
channel.exchangeDeclare(exchangeName, "topic", true, false, null);
channel.queueDeclare(queueName, true, false, false, null);
channel.queueBind(queueName, exchangeName, routingKey);
//1 限流方式 第一件事就是 autoAck设置为 false
//设置为1,表示一条一条数据处理
channel.basicQos(0, 1, false);
channel.basicConsume(queueName, false, new MyConsumer(channel));
}
}
复制代码
/** * * @ClassName: MyConsumer * @Description: TODO * @author Coder编程 * @date 2019年7月30日 下午23:11:55 * */
public class MyConsumer extends DefaultConsumer {
private Channel channel ;
public MyConsumer(Channel channel) {
super(channel);
this.channel = channel;
}
@Override
public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
System.err.println("-----------consume message----------");
System.err.println("consumerTag: " + consumerTag);
System.err.println("envelope: " + envelope);
System.err.println("properties: " + properties);
System.err.println("body: " + new String(body));
//须要作签收,false表示不支持批量签收
channel.basicAck(envelope.getDeliveryTag(), false);
}
}
复制代码
咱们先注释掉:channel.basicAck(envelope.getDeliveryTag(), false);而后启动Consumer。 查看Exchange微信
查看Queues 并发
而后再启动Producer。查看打印结果: ide
咱们会发现消费端,只收到了一条消息。这是为何呢?
第一点由于咱们在consumer中
channel.basicConsume(queueName, false, new MyConsumer(channel));
复制代码
第二个参数设置为false为手动签收。
第二点在qos中设置只接受一条消息。若是这一条消息不给Broker Ack应答的话,那么Broker会认为你并无消费完这一条消息,那么就不会继续发送消息。
channel.basicQos(0, 1, false);
复制代码
能够看下管控台,unack=1,Ready=4,total=5.
接下来咱们放开注释channel.basicAck(envelope.getDeliveryTag(), false); 进行消息签收。重启服务。
消费端进行消费的时候,若是因为业务异常咱们能够进行日志的记录,而后进行补偿!
若是因为服务器宕机等严重问题,那咱们就须要手工进行ACK保障消费端消费成功!
消费端重回队列是为了对没有处理成功的消息,把消息从新传递给Broker!
通常咱们在实际应用中,都会关闭重回队列,也就是设置为False.
/** * * @ClassName: Producer * @Description: 生产者 * @author Coder编程 * @date2019年7月30日 下午23:15:51 * */
public class Producer {
public static void main(String[] args) throws Exception {
//1建立ConnectionFactory
Connection connection = ConnectionUtils.getConnection();
Channel channel = connection.createChannel();
String exchange = "test_ack_exchange";
String routingKey = "ack.save";
for(int i =0; i<5; i ++){
Map<String, Object> headers = new HashMap<String, Object>();
headers.put("num", i);
//添加属性,后续会使用到
AMQP.BasicProperties properties = new AMQP.BasicProperties.Builder()
.deliveryMode(2) //投递模式,持久化
.contentEncoding("UTF-8")
.headers(headers)
.build();
String msg = "Hello RabbitMQ ACK Message " + i;
channel.basicPublish(exchange, routingKey, true, properties, msg.getBytes());
}
}
}
复制代码
/** * * @ClassName: Consumer * @Description: 消费者 * @author Coder编程 * @date2019年7月30日 下午23:13:51 * */
public class Consumer {
public static void main(String[] args) throws Exception {
//1建立ConnectionFactory
Connection connection = ConnectionUtils.getConnection();
Channel channel = connection.createChannel();
String exchangeName = "test_ack_exchange";
String queueName = "test_ack_queue";
String routingKey = "ack.#";
channel.exchangeDeclare(exchangeName, "topic", true, false, null);
channel.queueDeclare(queueName, true, false, false, null);
channel.queueBind(queueName, exchangeName, routingKey);
// 手工签收 必需要关闭 autoAck = false
channel.basicConsume(queueName, false, new MyConsumer(channel));
}
}
复制代码
/** * * @ClassName: MyConsumer * @Description: TODO * @author Coder编程 * @date 2019年7月30日 下午23:11:55 * */
public class MyConsumer extends DefaultConsumer {
private Channel channel ;
public MyConsumer(Channel channel) {
super(channel);
this.channel = channel;
}
@Override
public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
System.err.println("-----------consume message----------");
System.err.println("body: " + new String(body));
try {
Thread.sleep(2000);
} catch (InterruptedException e) {
e.printStackTrace();
}
if((Integer)properties.getHeaders().get("num") == 0) {
//Nack三个参数 第二个参数:是不是批量,第三个参数:是否重回队列(须要注意可能会发生重复消费,形成死循环)
channel.basicNack(envelope.getDeliveryTag(), false, true);
} else {
channel.basicAck(envelope.getDeliveryTag(), false);
}
}
}
复制代码
注意: 能够看到重回队列会出现重复消费致使死循环的问题,这时候最好设置重试次数,好比超过三次后,消息仍是消费失败,就将消息丢弃。
x-max-length 队列的最大大小 x-message-ttl 设置10秒钟,若是消息尚未被消费的话,就会被清除。
点击 test_ttl_exchange 进行绑定
生产端设置过时时间
AMQP.BasicProperties properties = new AMQP.BasicProperties.Builder()
.deliveryMode(2)
.contentEncoding("UTF-8")
.expiration("10000")
.headers(headers)
.build();
复制代码
这两个属性并不相同,一个对应的是消息体,一个对应的是队列的过时。
死信队列:DLX,Dead-Letter-Exchange RabbitMQ的死信队里与Exchange息息相关
消息变成死信有如下几种状况
DLX也是一个正常的Exchange,和通常的Exchange没有区别,它能在任何的队列上被指定,实际上就是设置某个队列的属性
当这个队列中有死信时,RabbitMQ就会自动的将这个消息从新发布到设置的Exchange上去,进而被路由到另外一个队列。
能够监听这个队列中消息作相应的处理,这个特征能够弥补RabbitMQ3.0之前支持的immediate参数的功能。
/** * * @ClassName: Producer * @Description: 生产者 * @author Coder编程 * @date2019年7月30日 下午23:15:51 * */
public class Producer {
public static void main(String[] args) throws Exception {
//建立ConnectionFactory
Connection connection = ConnectionUtils.getConnection();
Channel channel = connection.createChannel();
String exchange = "test_dlx_exchange";
String routingKey = "dlx.save";
String msg = "Hello RabbitMQ DLX Message";
for(int i =0; i<1; i ++){
AMQP.BasicProperties properties = new AMQP.BasicProperties.Builder()
.deliveryMode(2)
.contentEncoding("UTF-8")
.expiration("10000")
.build();
channel.basicPublish(exchange, routingKey, true, properties, msg.getBytes());
}
}
}
复制代码
/** * * @ClassName: Consumer * @Description: 消费者 * @author Coder编程 * @date2019年7月30日 下午23:13:51 * */
public class Consumer {
public static void main(String[] args) throws Exception {
//建立ConnectionFactory
Connection connection = ConnectionUtils.getConnection();
Channel channel = connection.createChannel();
// 这就是一个普通的交换机 和 队列 以及路由
String exchangeName = "test_dlx_exchange";
String routingKey = "dlx.#";
String queueName = "test_dlx_queue";
channel.exchangeDeclare(exchangeName, "topic", true, false, null);
Map<String, Object> agruments = new HashMap<String, Object>();
agruments.put("x-dead-letter-exchange", "dlx.exchange");
//这个agruments属性,要设置到声明队列上
channel.queueDeclare(queueName, true, false, false, agruments);
channel.queueBind(queueName, exchangeName, routingKey);
//要进行死信队列的声明:
channel.exchangeDeclare("dlx.exchange", "topic", true, false, null);
channel.queueDeclare("dlx.queue", true, false, false, null);
channel.queueBind("dlx.queue", "dlx.exchange", "#");
channel.basicConsume(queueName, true, new MyConsumer(channel));
}
}
复制代码
/** * * @ClassName: MyConsumer * @Description: TODO * @author Coder编程 * @date 2019年7月30日 下午23:11:55 * */
public class MyConsumer extends DefaultConsumer {
public MyConsumer(Channel channel) {
super(channel);
}
@Override
public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
System.err.println("-----------consume message----------");
System.err.println("consumerTag: " + consumerTag);
System.err.println("envelope: " + envelope);
System.err.println("properties: " + properties);
System.err.println("body: " + new String(body));
}
}
复制代码
运行Consumer,查看管控台 查看Exchanges
关闭Consumer,只运行Producer
过10秒钟后,消息过时
在咱们工做中,死信队列很是重要,用于消息没有消费者,处于死信状态。咱们能够才用补偿机制。
本次主要介绍了RabbitMQ的高级特性,首先介绍了互联网大厂在实际使用中如何保障100%的消息投递成功和幂等性的,以及对RabbitMQ的确认消息、返回消息、ACK与重回队列、消息的限流,以及对超时时间、死信队列的使用
欢迎关注我的微信公众号:Coder编程 获取最新原创技术文章和免费学习资料,更有大量精品思惟导图、面试资料、PMP备考资料等你来领,方便你随时随地学习技术知识! 新建了一个qq群:315211365,欢迎你们进群交流一块儿学习。谢谢了!也能够介绍给身边有须要的朋友。
文章收录至 Github: github.com/CoderMerlin… Gitee: gitee.com/573059382/c… 欢迎关注并star~
![]()
参考文章:
《RabbitMQ消息中间件精讲》
推荐文章:
消息中间件——RabbitMQ(五)快速入门生产者与消费者,SpringBoot整合RabbitMQ!