spring boot 1.4连接rabbitmq,并使用延时队列插件
·
因为spring boot 1.4版本比较老了,所以不能直接支持延时队列。所以需要用手工创建Exchange
@RabbitListener(
queues = "test-log-queue"
)
public void processMessage(String message) {
logger.info("我接收了[{}]。", message);
}
@Bean
public ApplicationRunner runner(@Qualifier("amqpTemplate") RabbitTemplate template) {
return args -> {
MessageProperties messageProperties = new MessageProperties();
template.convertAndSend("test-log-exchange", "main", "34020000001320000001", message->{
message.getMessageProperties().setHeader("x-delay", 60 * 1000);
return message;
});
logger.info("我发送成功了!");
};
}
@Bean
public org.springframework.amqp.core.Queue manualDelayQueue() {
// 持久化、非排他、非自动删除
return new org.springframework.amqp.core.Queue("test-log-queue", true, false, false);
}
// 3. 手动声明交换机 + 绑定队列(核心:PostConstruct初始化)
@PostConstruct
public void declareManualExchangeAndBinding() {
// 3.1 构建自定义延迟交换机(x-delayed-message类型)
Map<String, Object> exchangeArgs = new HashMap<>();
exchangeArgs.put("x-delayed-type", "direct"); // 指定延迟交换机的路由类型
Map<String, Object> args = new HashMap<>();
args.put("x-delayed-type", "direct"); // 指定内部交换机类型为 direct
org.springframework.amqp.core.Exchange manualDelayExchange = new CustomExchange("test-log-exchange", "x-delayed-message", true, false, args);
// 3.2 手动创建交换机(绕过Spring AMQP的类型校验)
amqpAdmin.declareExchange(manualDelayExchange);
// 3.3 手动绑定队列到交换机(指定路由键)
Binding binding = BindingBuilder.bind(manualDelayQueue())
.to(manualDelayExchange)
.with("main")
.noargs();
amqpAdmin.declareBinding(binding);
logger.info("手动声明交换机&绑定完成:" + "test-log-exchange");
}
更多推荐



所有评论(0)