rabbitmq消息队列学习之进阶
设置队列的TTL
通过channel.queueDeclare方法中的x-expires参数可以控制队列被自动删除前处于未使用状态的时间。未使用的意思是队列傻姑娘没有任何消费者,队列也没有被重新声明,并且在过期时间段内也未调用过Basic.Get命令。
用于表示过期时间的 x-expires参数 以毫秒 为单位,并且服从和 x-message-ttl 一样的约束条件,不过不能设置为0,比如该参数设置为 1000,则表示该队列如果在1秒钟之内为使用则会被删除。
package com.study.rabbit.ttl;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.MessageProperties;
import java.io.IOException;
import java.util.HashMap;
import java.util.Map;
import java.util.concurrent.TimeoutException;
public class DirectProducer {
public static void main(String[] args) throws IOException, TimeoutException {
Connection connection = null;
Channel channel = null;
try {
//通过连接工厂创建新的连接和mq进行连接
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("127.0.0.1"); //设置mq所在服务器的ip
factory.setPort(5672); //设置mq所在服务器的端口号
factory.setUsername("guest");//设置登录用户名
factory.setPassword("guest");//设置登录密码
factory.setVirtualHost("/");//rabbitmq默认虚拟机名称为“/”,相当于一个独立的mq服务
//创建与RabbitMQ服务的TCP连接
connection = factory.newConnection();
//创建与Exchange(交换机)的通道,每个连接可以创建多个通道,每个通道代表一个会话任务
channel = connection.createChannel();
channel.exchangeDeclare("exchange.normal", "direct", true);
Map<String, Object> argMap = new HashMap<>();
argMap.put("x-expires", 10000);
channel.queueDeclare("queue.normal", true, false, false, argMap);
channel.queueBind("queue.normal", "exchange.normal", "expires.key");
channel.basicPublish("exchange.normal", "expires.key", MessageProperties.PERSISTENT_TEXT_PLAIN, "dlx".getBytes());
} catch (Exception e) {
e.printStackTrace();
} finally {
if (channel != null) {
channel.close();
}
if (connection != null) {
connection.close();
}
}
}
}

死信队列
DLX,全称为Dead-Letter-Exchange,可以称之为私信交换器,也有人称之为死信邮箱。当消息在一个队列中变为死信(dead message) 之后,它能被重新发送到另一个交换器中,这个交换器就是DLX,绑定DLX的队列就称之为死信队列。
消息变为死信一般是由于以下几种情况:
- 消息被拒绝( Basic.Reject/Basic.Nack ),并且设置requeue 参数为false;
- 消息过期
- 队列到达最大长度
DLX 也是一个正常的交换器,和一般的交换器没有区别,它能在任何的队列上被指定,实际上就是设置某个队列的属性。当这个队列中存在死信时,RabbitMQ 就会自动地将这个消息重新发布到设置的DLX上去,进而被路由到另一个队列,即死信队列。可以监听这个队列中的消息已进行相应的处理,这个特性与将消息的TTL设置为0配合使用可以弥补immediate参数的功能。

生产者代码
package com.study.rabbit.dlx;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.MessageProperties;
import java.io.IOException;
import java.util.HashMap;
import java.util.Map;
import java.util.concurrent.TimeoutException;
public class DirectProducer {
public static void main(String[] args) throws IOException, TimeoutException {
Connection connection = null;
Channel channel = null;
try {
//通过连接工厂创建新的连接和mq进行连接
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("127.0.0.1"); //设置mq所在服务器的ip
factory.setPort(5672); //设置mq所在服务器的端口号
factory.setUsername("guest");//设置登录用户名
factory.setPassword("guest");//设置登录密码
factory.setVirtualHost("/");//rabbitmq默认虚拟机名称为“/”,相当于一个独立的mq服务
//创建与RabbitMQ服务的TCP连接
connection = factory.newConnection();
//创建与Exchange(交换机)的通道,每个连接可以创建多个通道,每个通道代表一个会话任务
channel = connection.createChannel();
channel.exchangeDeclare("exchange.dlx", "direct", true);
channel.exchangeDeclare("exchange.normal", "fanout", true);
Map<String, Object> argMap = new HashMap<>();
argMap.put("x-message-ttl", 10000);
argMap.put("x-dead-letter-exchange", "exchange.dlx");
argMap.put("x-dead-letter-routing-key", "routingkey");
channel.queueDeclare("queue.normal", true, false, false, argMap);
channel.queueBind("queue.normal", "exchange.normal", "");
channel.queueDeclare("queue.dlx", true, false, false, null);
channel.queueBind("queue.dlx", "exchange.dlx", "routingkey");
channel.basicPublish("exchange.normal", "rk", MessageProperties.PERSISTENT_TEXT_PLAIN, "dlx".getBytes());
} catch (Exception e) {
e.printStackTrace();
} finally {
if (channel != null) {
channel.close();
}
if (connection != null) {
connection.close();
}
}
}
}
http://localhost:15672/#/queues

延迟队列
延迟队列存储的对象是对应的延迟消息,所谓“延迟消息”是指当消息被发送以后,并不想让消费者立刻拿到消息,而是等待特定时间后,消费者才能拿到消息进行消费。
延迟队列的使用场景很多,比如:
在订单系统中,一个用户下单之后通常有30分钟的时间进行支付,如果30分钟之内没有支付成功,那么这个订单将进行取消,这时就可以使用延迟队列来处理这些订单了。

备份交换器
备份交换器,如果即不想复杂化生产者的编程逻辑,又不想消息丢失,那么可以使用备份交换器,这样可以将未被路由的消息存储在RabbitMQ ,再在需要的时候去处理这些消息。

生产者确认
在使用RabbitMQ 的时候,可以通过消息持久化操作来解决因为服务器的异常崩溃而导致的消息丢失。除此之外,我们还会遇到一个问题,当消息的生产者将消息发送出去之后,消息到底有没有正确到达服务器呢?如果不进行特殊配置,默认情况下发送消息的操作是不会返回任何信息给生产者的,也就是默认情况下生产者是不知道消息有没有正确到达服务器,如果在消息到达服务器之前已经丢失,持久化操作也解决不来这个问题,因为消息根本没有到达服务器,何谈持久化。
RabbitMQ 针对这个问题,提供了两种解决方案:
-
事务机制

-
发送方确认机制

package com.study.rabbit.confirm;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.MessageProperties;
import java.io.IOException;
import java.util.concurrent.TimeoutException;
public class DirectProducer {
public static void main(String[] args) throws IOException, TimeoutException {
Connection connection = null;
Channel channel = null;
try {
//通过连接工厂创建新的连接和mq进行连接
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("127.0.0.1"); //设置mq所在服务器的ip
factory.setPort(5672); //设置mq所在服务器的端口号
factory.setUsername("guest");//设置登录用户名
factory.setPassword("guest");//设置登录密码
factory.setVirtualHost("/");//rabbitmq默认虚拟机名称为“/”,相当于一个独立的mq服务
//创建与RabbitMQ服务的TCP连接
connection = factory.newConnection();
//创建与Exchange(交换机)的通道,每个连接可以创建多个通道,每个通道代表一个会话任务
channel = connection.createChannel();
channel.exchangeDeclare("exchange.normal", "direct", true);
channel.queueDeclare("queue.normal", true, false, false, null);
channel.queueBind("queue.normal", "exchange.normal", "expires.key");
try {
// 将信道设置为 publisher confirm 模式
channel.confirmSelect();
// 之后正常发送消息
channel.basicPublish("exchange.normal", "expires.key", MessageProperties.PERSISTENT_TEXT_PLAIN, "dlx".getBytes());
if(!channel.waitForConfirms()){
System.out.println("send message failed ");
// do something
}
}catch (Exception ex){
}
} catch (Exception e) {
e.printStackTrace();
} finally {
if (channel != null) {
channel.close();
}
if (connection != null) {
connection.close();
}
}
}
}

更多推荐



所有评论(0)