死信队列
一 、队列的创建和绑定
二、消费者消费,消费失败之后的逻辑
@Component
@Log4j2
@RabbitListener(queues = REAL_DEAD_QUEUE)
public class PayOrderConsumer {
@Resource
OrderMapper orderMapper;
@Resource
RabbitTemplate rabbitTemplate;
@RabbitHandler
public void receiveDeadMsg(String msg, Channel channel, Message message) throws IOException {
JSONObject msgJson = JSON.parseObject(msg);
int count = Integer.parseInt(msgJson.getString("count"));
Order order=null;
try {
log.info("接收到消息:{}",msg);
//业务逻辑
....
channel.basicAck(message.getMessageProperties().getDeliveryTag(), Boolean.FALSE);
}catch (Exception e){
if (count>2){
log.info("重试次数超过三次失败的消息:{}",msg)
//发生异常三次重试都失败的时候的业务逻辑
....
channel.basicAck(message.getMessageProperties().getDeliveryTag(), Boolean.FALSE);
throw e;
}else {
//消费消息失败的时候的逻辑
log.error("消息{}即将再次返回队列处理...",msg);
msgJson.put("count",count+1); //计数器
rabbitTemplate.convertAndSend(NORMAL_EXCHANGE, NORMAL_ROUTING_KEY, msgJson.toString(), new MessagePostProcessor() {
@Override
public Message postProcessMessage(Message message) throws AmqpException {
MessageProperties messageProperties = message.getMessageProperties();
//设置过期时间TTL
messageProperties.setExpiration(String.valueOf(EXPIRE));
return message;
}
});
}
channel.basicAck(message.getMessageProperties().getDeliveryTag(), Boolean.FALSE);
}
}
}
三、发送消息
JSONObject jsonObject = new JSONObject();
jsonObject.put("msg",orderNumber);
jsonObject.put("count",0);
rabbitTemplate.convertAndSend(PayOrderRabbitConfig.NORMAL_EXCHANGE,PayOrderRabbitConfig.NORMAL_ROUTING_KEY,jsonObject);
log.info("生成订单{}成功",orderNumber);