rabbitMQConfig
@Configuration
public class RabbitMQConfig {
public static final String TOPIC_EXCHANGE_NAME = "exchange_topic_boot";
public static final String TOPIC_QUEUE_NAME_1 = "queue_topic_boot1";
//创建一个交换机
@Bean("topicExchange")
public Exchange topicExchange(){
return ExchangeBuilder.topicExchange(TOPIC_EXCHANGE_NAME).durable(true).build();
}
//创建一个队列
@Bean("topicQueue")
public Queue topicQueue(){
return QueueBuilder.durable(TOPIC_QUEUE_NAME_1).build();
}
//绑定交换机和队列
@Bean
public Binding basicBinding(@Qualifier("topicExchange") Exchange exchange,@Qualifier("topicQueue") Queue queue){
return BindingBuilder.bind(queue).to(exchange).with("*.peng").noargs();
}
}
Consumer实现:通过监听器的方式
@Service
public class ConsumerListener {
@RabbitListener(queues = RabbitMQConfig.TOPIC_QUEUE_NAME_1)
public void QueueListener(Message msg){
System.out.println("consumer:"+msg);
}
}
rabbitMQ高级特性#
1、确认消费机制
从生产者发送信息到消费者,为了确保信息发送成功、信息不会丢失和确认消费。
RabbitMQ提供了以下解决方法:
- 信息发送成功:生产者发送消息到exchange交换机,exchange转发数据到队列,这两个步骤都可能出现消息发送失败的可能。解决方案如下:
@Bean
public RabbitTemplate rabbitTemplate(){
connectionFactory.setPublisherConfirmType(CachingConnectionFactory.ConfirmType.CORRELATED);
connectionFactory.setPublisherReturns(true);
RabbitTemplate rabbitTemplate = new RabbitTemplate(connectionFactory);
//检测信息能否成功的发送到交换机
rabbitTemplate.setConfirmCallback(new RabbitTemplate.ConfirmCallback() {
// ack:当ack为false时数据发送失败,result为失败的原因
@Override
public void confirm(CorrelationData correlationData, boolean ack, String result) {
if(!ack){
System.out.println("数据未能发送到交换机");
System.out.println("原因是:"+result);
}else
{
System.out.println("数据成功发送到交换机");
}
}
});
//发送完消息后输出反馈信息
rabbitTemplate.setMandatory(true);
rabbitTemplate.setReturnsCallback(new RabbitTemplate.ReturnsCallback() {
@Override
public void returnedMessage(ReturnedMessage returnedMessage) {
System.out.println("exchange转发数据到路由失败,执行回退操作");
System.out.println(returnedMessage.toString());
}
});
return rabbitTemplate;
}
当我发送消息到不存在的交换机时,消息发送到交换机失败,控制台输出如下:
当交换机使用不存在的路由转发消息时,出现失败,执行回退操作:
-
信息不会丢失:信息丢失指的是在RabbitMQ服务器宕机或者崩溃时,队列中大量的消息还没有被消费,重启后消息都丢失了,可以在创建队列和交换机时设置持久化参数为true。
-
消息确认消费(Consumer Acknowledge):rabbitmq的消息确认消费机制有三种,分别是none、auto、manual,其中默认为none,比较有用的的是manual,也就是手动设置。
创建消费者监听器工厂,设置模式为manual。
@Bean(SINGLE_LISTENER_FACTORY)
public SimpleRabbitListenerContainerFactory singleListener(){
SimpleRabbitListenerContainerFactory factory = new SimpleRabbitListenerContainerFactory();
factory.setConnectionFactory(connectionFactory);
factory.setAcknowledgeMode(AcknowledgeMode.MANUAL);
factory.setConcurrentConsumers(1);
factory.setMaxConcurrentConsumers(1);
return factory;
}
创建监听器
@RabbitListener(queues = RabbitMQConfig.TOPIC_QUEUE_NAME_1,
containerFactory = RabbitMQConfig.SINGLE_LISTENER_FACTORY)
public void QueueListener(Message msg, Channel channel) throws IOException, InterruptedException {
// Thread.sleep(10*1000);
long deliveryTag = msg.getMessageProperties().getDeliveryTag();
try {
//这里可以写消费者对消息的处理操作
System.out.println("消息处理逻辑");
int i = 1/0;
channel.basicAck(deliveryTag,false);
} catch (Exception e) {
System.out.println("消息未被消费,加入死信队列中");
channel.basicNack(deliveryTag,false,false);
}
}
死信队列
消息变为死信一般有三种情况:
1、消息被消费者拒绝
2、消息在TTL时间内没有被消费。其中TTL分为消息的TTL和队列的TTL。当一个消息具有TTL并且在一个由TTL的队列时,会按照较小的TTL设置过期时间。
3、队列满了
默认情况下死信会丢死,为了避免死信的丢失我们可以设置死信交换机和死信队列来管理死信。
- 可以通过参数设置队列的TTL。
//设置队列为延迟队列,过期时间为3秒,并且绑定死信交换机DEATH_EXCHANGE, 当消息过期时会被发送到死信交换机,随后由死信交换机处理。
@Bean(TOPIC_QUEUE_NAME_1)
public Queue topicQueue(){
Map args = new HashMap<>(2);
args.put("x-dead-letter-exchange", DEATH_EXCHANGE);
args.put("x-dead-letter-routing-key", "death");
args.put("x-message-ttl", 3*1000);
Queue queue = QueueBuilder.durable(TOPIC_QUEUE_NAME_1).withArguments(args).build();
return queue;
}
- 在发送消息时可以指定消息的TTL(redis中可以通过expire命令设置key的过期时间,还挺像的)
rabbitTemplate.convertAndSend(RabbitMQConfig.TOPIC_EXCHANGE_NAME, "death", "你好呀", new MessagePostProcessor() {
@Override
public Message postProcessMessage(Message message) throws AmqpException {
//设置消息的过期时间
message.getMessageProperties().setExpiration("5000");
return message;
}
});
消费端限流
在消费者监听器工厂设置以下参数。
factory.setPrefetchCount(n);
当限流量为n时,消费者每次会向队列中拉取n个消息进行处理,而其他的消息则会排队,可以限制短时间内大量消息的访问。我们可以根据系统的并发承受压力来设置这个值。
延迟队列是死信队列的一种应用,为一个队列设置过期时间(TTL),但不设置消费者,同时绑定一个死信交换机,当消息过期后会被发送到死信交换机,然后再由死信交换机绑定队列,队列发送消息给消费者,这样就实现了延迟队列。



