springboot rabbitmq 延时队列
配置:
spring:
rabbitmq:
host: 127.0.0.1
port: 5672
username: guest
password: guest
virtual-host: /
publisher-confirms: true
config:
controller:
@PostMapping("/feign/createOrder")
public void createOrder(@RequestBody OrderTradeRecordRequest req) {
testService.createOrder(req);
devService.send(req);
System.out.println("===============sender已发送====================="+req);
}
sender:
@Component
public class DevService {
@Autowired
private AmqpTemplate amqpTemplate;
//有两种方法
public void send(OrderTradeRecordRequest req){
System.out.println("【订单生成时间】" + new Date().toString() +"【1分钟后检查订单是否已经支付】" + req.toString() );
//第一种,使用了jdk8的新特性
/*this.amqpTemplate.convertAndSend(DelayRabbitConfig.ORDER_DELAY_EXCHANGE, DelayRabbitConfig.ORDER_DELAY_ROUTING_KEY, req, message -> {
// 如果配置了 params.put("x-message-ttl", 5 * 1000); 那么这一句也可以省略,具体根据业务需要是声明 Queue 的时候就指定好延迟时间还是在发送自己控制时间
message.getMessageProperties().setExpiration(1 * 1000 * 60 + "");
return message;
});*/
//this.amqpTemplate.convertAndSend(DelayRabbitConfig.ORDER_QUEUE_NAME, req);
//第二种
MessagePostProcessor msg = new MessagePostProcessor() {
@Override
public Message postProcessMessage(Message msg) throws AmqpException {
// 设置延迟毫秒值
msg.getMessageProperties().setExpiration("5000");
return msg;
}
};
this.amqpTemplate.convertAndSend(DelayRabbitConfig.ORDER_DELAY_EXCHANGE, DelayRabbitConfig.ORDER_DELAY_ROUTING_KEY, req, msg);
}
}
customer:
@Component
public class RabbitMQListener {
@Autowired
private TestService testService;
//直接消费模式
@RabbitListener(queues = {DelayRabbitConfig.ORDER_QUEUE_NAME})
public void consumeMessage(OrderTradeRecordRequest req){
System.out.println("===============接收==================="+req);
//TODO:表示时间已经到了,却还没付款,则需要处理为失效
Integer status = req.getStatus();
if (status.equals(1)||status==1){
Integer orderId=req.getOrderId();
testService.updateOrder(orderId);
}
}
}
