本地消息表 + 定时任务
本地消息表:主要用于存储 业务数据、交换机、队列、路由、次数
定时任务:定时扫描本地消息表,重新给业务队列投递消息。
利用 rabbitmq_delayed_message_exchange 插件 实现延迟队列
具体思路:业务队列消费失败时,给延迟队列发送一条消息,消息包含业务数据、交换机、队列、次数、最大次数等,延迟队列收到消息后重新给业务队列投递消息。业务队列二次收到消息时,再次消费失败,校验最大次数,判断是否再次重试。
- pom.xml
run.siyuan siyuan-common 1.0-SNAPSHOT org.springframework.boot spring-boot-starter-web org.springframework.boot spring-boot-starter-amqp org.projectlombok lombok
- application.yml
server:
port: 8080
spring:
rabbitmq:
addresses: 127.0.0.1
port: 5672
username: siyuan
password: siyuan123456
virtual-host: /
- PluginDelayRabbitConfig.java
import com.rabbitmq.client.ConnectionFactory;
import org.springframework.amqp.core.*;
import org.springframework.amqp.rabbit.core.RabbitAdmin;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.amqp.support.converter.Jackson2JsonMessageConverter;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import java.util.HashMap;
import java.util.Map;
@Configuration
public class PluginDelayRabbitConfig {
@Bean("pluginDelayExchange")
public CustomExchange pluginDelayExchange() {
Map argMap = new HashMap<>();
argMap.put("x-delayed-type", "direct");//必须要配置这个类型,可以是direct,topic和fanout
//第二个参数必须为x-delayed-message
return new CustomExchange("PLUGIN_DELAY_EXCHANGE","x-delayed-message",false, false, argMap);
}
@Bean("pluginDelayQueue")
public Queue pluginDelayQueue(){
return new Queue("PLUGIN_DELAY_QUEUE");
}
@Bean
public Binding pluginDelayBinding(@Qualifier("pluginDelayQueue") Queue queue, @Qualifier("pluginDelayExchange") CustomExchange customExchange){
return BindingBuilder.bind(queue).to(customExchange).with("delay").noargs();
}
}
- RabbitmqConsumer.java
import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.rabbit.annotation.RabbitHandler;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component;
import java.text.SimpleDateFormat;
import java.util.Date;
@Slf4j
@Component
public class RabbitmqConsumer {
@RabbitHandler
@RabbitListener(queues = "PLUGIN_DELAY_QUEUE")//监听延时队列
public void fanoutConsumer(String msg){
SimpleDateFormat sdf = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss");
System.out.println("【插件延迟队列】【" + sdf.format(new Date()) + "】收到消息:" + msg);
}
}
- RabbitMqController.java
import cn.hutool.json.JSONObject;
import cn.hutool.json.JSONUtil;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.core.MessageProperties;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RequestParam;
import org.springframework.web.bind.annotation.RestController;
import java.text.SimpleDateFormat;
import java.util.Date;
@RestController
public class RabbitMqController {
@Autowired
RabbitTemplate rabbitTemplate;
@GetMapping(value = "/plugin/send")
public String pluginMsgSend(@RequestParam Integer time) {
JSONObject json = new JSONObject();
json.set("name", "插件延迟消息");
json.set("time", System.currentTimeMillis());
json.set("delayTime", time);
MessageProperties messageProperties = new MessageProperties();
messageProperties.setHeader("x-delay", 1000 * time);//延迟5秒被删除
Message message = new Message(JSONUtil.toJsonStr(json).getBytes(), messageProperties);
rabbitTemplate.convertAndSend("PLUGIN_DELAY_EXCHANGE", "delay", message);//交换机和路由键必须和配置文件类中保持一致
SimpleDateFormat sdf = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss");
System.out.println("消息发送成功【" + sdf.format(new Date()) + "】" + "延迟时间:" + time);
return "succ";
}
}
方案三:
利用 TTL 消息 + DLX 死信队列 实现延迟队列
具体思路:业务队列消费失败时,会发送一条TTL 消息,消息包含业务数据、交换机、队列、次数、最大次数等,TTL 消息过期后会进入死信队列,此时监听死信队列接收消息,校验是否达到重试次数,再重新投递给业务队列,业务队列二次收到消息时,再次消费失败,校验最大次数,判断是否再次重试。超过最大次数入库,人工干预处理
- pom.xml
run.siyuan siyuan-common 1.0-SNAPSHOT org.springframework.boot spring-boot-starter-web org.springframework.boot spring-boot-starter-amqp org.projectlombok lombok
- application.yml
server:
port: 8080
spring:
rabbitmq:
addresses: 127.0.0.1
port: 5672
username: siyuan
password: siyuan123456
virtual-host: /
- Constants.java
public interface Constants {
// ------------------------------ delay -------------------------------------
// 延时交换机
String DELAY_EXCHANGE = "delay.exchange";
// 延时交换机队列
String DELAY_EXCHANGE_QUEUE = "delay.exchange.queue";
// 延时交换机路由键
String DELAY_EXCHANGE_ROUTE_KEY = "delay.exchange.route.key";
// ------------------------------ dead.letter.fanout -------------------------------------
// 死信交换机
String DELAY_LETTER_EXCHANGE = "dead.letter.exchange";
// 死信交换机队列
String DELAY_LETTER_EXCHANGE_QUEUE = "dead.letter.exchange.queue";
// 死信交换机路由键
String DELAY_LETTER_EXCHANGE_ROUTE_KEY = "dead.letter.exchange.route.key";
// ------------------------------ 业务队列 -------------------------------------
String SERVICE_QUEUE = "service.queue";
}
- RetryRabbitConfig.java
import org.springframework.amqp.core.*;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import run.siyuan.common.rabbitmq.Constants;
import java.util.HashMap;
import java.util.Map;
@Configuration
public class RetryRabbitConfig {
@Bean
public DirectExchange ttlDelayExchangeRetry() {
return new DirectExchange(Constants.DELAY_EXCHANGE);
}
@Bean
public Queue ttlDelayExchangeQueueRetry() {
Map map = new HashMap();
//队列中所有消息5秒后过期
//map.put("x-message-ttl", 1000 * 60 * 5);
//过期后进入死信队列
map.put("x-dead-letter-exchange", Constants.DELAY_LETTER_EXCHANGE);
return new Queue(Constants.DELAY_EXCHANGE_QUEUE, false, false, false, map);
}
@Bean
public Binding bindTtlExchangeAndQueueRetry() {
return BindingBuilder.bind(ttlDelayExchangeQueueRetry()).to(ttlDelayExchangeRetry()).with(Constants.DELAY_EXCHANGE_ROUTE_KEY);
}
@Bean
public FanoutExchange deadLetterExchange() {
return new FanoutExchange(Constants.DELAY_LETTER_EXCHANGE);
}
@Bean
public Queue deadLetterQueue() {
return new Queue(Constants.DELAY_LETTER_EXCHANGE_QUEUE);
}
@Bean
public Queue serviceQueue() {
return new Queue(Constants.SERVICE_QUEUE);
}
@Bean
public Binding deadLetterBind() {
return BindingBuilder.bind(deadLetterQueue()).to(deadLetterExchange());
}
}
- MessageRetryVo
package run.siyuan.rabbitmq.retry.message.model;
import lombok.Data;
import lombok.EqualsAndHashCode;
import lombok.experimental.Accessors;
import java.io.Serializable;
import java.util.Date;
@Data
@EqualsAndHashCode(callSuper = false)
@Accessors(chain = true)
public class MessageRetryVo implements Serializable {
private static final long serialVersionUID = 1L;
private String bodyMsg;
private String exchangeName;
private String routingKey;
private String queueName;
private Integer maxTryCount = 3;
private Integer currentRetryCount = 0;
private String errorMsg;
private Date createTime;
private Integer type;
@Override
public String toString() {
return "MessageRetryDTO{" +
"bodyMsg='" + bodyMsg + ''' +
", exchangeName='" + exchangeName + ''' +
", routingKey='" + routingKey + ''' +
", queueName='" + queueName + ''' +
", maxTryCount=" + maxTryCount +
", currentRetryCount=" + currentRetryCount +
", errorMsg='" + errorMsg + ''' +
", createTime=" + createTime +
'}';
}
public boolean checkRetryCount(Integer type) {
//检查重试次数是否超过最大值
if (this.currentRetryCount <= this.maxTryCount) {
if (type.equals(0)) {
retryCountCalculate();
}
return true;
}
return false;
}
private void retryCountCalculate() {
this.currentRetryCount = this.currentRetryCount + 1;
}
}
- ServiceConsumer.java
import cn.hutool.json.JSONObject;
import cn.hutool.json.JSONUtil;
import com.rabbitmq.client.Channel;
import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.amqp.support.AmqpHeaders;
import org.springframework.messaging.handler.annotation.Header;
import org.springframework.stereotype.Component;
import run.siyuan.common.rabbitmq.Constants;
import run.siyuan.rabbitmq.retry.message.service.CommonMessageDelayService;
import java.io.IOException;
@Slf4j
@Component
public class ServiceConsumer extends CommonMessageDelayService {
@RabbitListener(queues = Constants.SERVICE_QUEUE, ackMode = "MANUAL", concurrency = "1")
private void consumer(Message message, @Header(AmqpHeaders.DELIVERY_TAG) long deliveryTag, Channel channel) throws IOException {
try {
byte[] body = message.getBody();
String msg = new String(body);
log.info("【正常队列】【" + System.currentTimeMillis() + "】收到死信队列消息:" + msg);
JSONObject json = JSONUtil.parseObj(msg);
if (json.getInt("id") < 0) {
throw new Exception("id 小于 0");
}
channel.basicAck(deliveryTag, false);
} catch (Exception e) {
log.info("消费异常:{}", e.getMessage());
channel.basicNack(deliveryTag, false, false);
sendDelayMessage(message, e);
}
}
}
- DeadLetterConsumer
import com.rabbitmq.client.Channel;
import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.rabbit.annotation.RabbitHandler;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.amqp.support.AmqpHeaders;
import org.springframework.messaging.handler.annotation.Header;
import org.springframework.stereotype.Component;
import run.siyuan.common.rabbitmq.Constants;
import run.siyuan.rabbitmq.retry.message.service.CommonMessageRetryService;
import java.io.IOException;
@Slf4j
@Component
public class DeadLetterConsumer extends CommonMessageRetryService {
@RabbitHandler
@RabbitListener(queues = Constants.DELAY_LETTER_EXCHANGE_QUEUE, ackMode = "MANUAL", concurrency = "1")
public void consumer(Message message, @Header(AmqpHeaders.DELIVERY_TAG) long deliveryTag, Channel channel) throws IOException {
try {
log.info("【死信队列】【" + System.currentTimeMillis() + "】收到死信队列消息:", new String(message.getBody()));
retryMessage(message);
channel.basicAck(deliveryTag, false);
} catch (Exception e) {
channel.basicNack(deliveryTag, false, false);
}
}
}
- CommonMessageDelayService.java
import cn.hutool.json.JSONUtil;
import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.core.MessageProperties;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import run.siyuan.common.rabbitmq.Constants;
import run.siyuan.rabbitmq.retry.message.model.MessageRetryVo;
@Slf4j
public abstract class CommonMessageDelayService extends AbstractCommonMessageService {
@Autowired
private RabbitTemplate rabbitTemplate;
protected void sendDelayMessage(Message message, Exception e) {
try {
//封装消息
MessageRetryVo delayMessageVo = buildMessageRetryInfo(message);
log.info("延时消息:{}", delayMessageVo);
//获取所有堆栈信息
StackTraceElement[] stackTraceElements = e.getStackTrace();
//默认的异常类全路径为第一条异常堆栈信息的
String exceptionClassTotalName = stackTraceElements[0].toString();
//遍历所有堆栈信息,找到vip.xiaonuo开头的第一条异常信息
for (StackTraceElement stackTraceElement : stackTraceElements) {
if (stackTraceElement.toString().contains("com.central")) {
exceptionClassTotalName = stackTraceElement.toString();
break;
}
}
log.info("异常信息:{}", exceptionClassTotalName);
delayMessageVo.setErrorMsg(exceptionClassTotalName);
delayMessageVo.setType(0);
prepareAction(delayMessageVo);
} catch (Exception exception) {
log.warn("处理消息异常,错误信息:", exception);
}
}
@Override
protected void sendMessage(MessageRetryVo retryVo) {
//将补偿消息实体放入头部,原始消息内容保持不变
MessageProperties messageProperties = new MessageProperties();
// 消息的有效时间固定,不使用自定义时间
messageProperties.setExpiration(String.valueOf(1000 * 10 * 1));
messageProperties.setHeader("message_retry_info", JSONUtil.toJsonStr(retryVo));
Message ttlMessage = new Message(JSONUtil.toJsonStr(retryVo).getBytes(), messageProperties);
rabbitTemplate.convertAndSend(Constants.DELAY_EXCHANGE, Constants.DELAY_EXCHANGE_ROUTE_KEY, ttlMessage);
log.info("发送业务消息 完成 时间:{}", System.currentTimeMillis());
}
}
- CommonMessageRetryService.java
import cn.hutool.json.JSONUtil;
import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.core.MessageProperties;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import run.siyuan.rabbitmq.retry.message.model.MessageRetryVo;
@Slf4j
public abstract class CommonMessageRetryService extends AbstractCommonMessageService {
@Autowired
private RabbitTemplate rabbitTemplate;
public void retryMessage(Message message) {
try {
//封装消息
MessageRetryVo retryMessageVo = buildMessageRetryInfo(message);
log.info("重试消息:{}", retryMessageVo);
retryMessageVo.setType(1);
prepareAction(retryMessageVo);
} catch (Exception exception) {
log.warn("处理消息异常,错误信息:", exception);
}
}
@Override
protected void sendMessage(MessageRetryVo retryVo) {
//将补偿消息实体放入头部,原始消息内容保持不变
MessageProperties messageProperties = new MessageProperties();
messageProperties.setHeader("message_retry_info", JSONUtil.toJsonStr(retryVo));
Message message = new Message(retryVo.getBodyMsg().getBytes(), messageProperties);
rabbitTemplate.convertAndSend(retryVo.getExchangeName(), retryVo.getRoutingKey(), message);
log.info("发送业务消息 完成 时间:{}", System.currentTimeMillis());
}
}
- AbstractCommonMessageService
import cn.hutool.json.JSONUtil;
import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.core.Message;
import run.siyuan.rabbitmq.retry.message.model.MessageRetryVo;
import java.nio.charset.StandardCharsets;
import java.util.Date;
import java.util.Map;
import java.util.Objects;
@Slf4j
public abstract class AbstractCommonMessageService {
protected abstract void sendMessage(MessageRetryVo retryVo);
protected MessageRetryVo buildMessageRetryInfo(Message message) {
//如果头部包含补偿消息实体,直接返回
Map messageHeaders = message.getMessageProperties().getHeaders();
if (messageHeaders.containsKey("message_retry_info")) {
Object retryMsg = messageHeaders.get("message_retry_info");
if (Objects.nonNull(retryMsg)) {
return JSONUtil.toBean(JSONUtil.parseObj(retryMsg), MessageRetryVo.class);
}
}
//自动将业务消息加入补偿实体
MessageRetryVo messageVo = new MessageRetryVo();
messageVo.setBodyMsg(new String(message.getBody(), StandardCharsets.UTF_8));
messageVo.setExchangeName(message.getMessageProperties().getReceivedExchange());
messageVo.setRoutingKey(message.getMessageProperties().getReceivedRoutingKey());
messageVo.setQueueName(message.getMessageProperties().getConsumerQueue());
messageVo.setCreateTime(new Date());
return messageVo;
}
protected void prepareAction(MessageRetryVo messageVo) {
if (messageVo.checkRetryCount(messageVo.getType())) {
this.sendMessage(messageVo);
} else {
if (log.isWarnEnabled()) {
log.warn("当前任务重试次数已经到达最大次数,业务数据:" + messageVo.toString());
}
doFailCallBack(messageVo);
}
}
protected void doFailCallBack(MessageRetryVo messageVo) {
try {
saveRetryMessageInfo(messageVo);
} catch (Exception e) {
log.warn("执行失败回调异常,错误原因:{}", e.getMessage());
}
}
protected void saveRetryMessageInfo(MessageRetryVo messageVo) {
try {
log.info("重试消息次数:{} message_retry_info:{}", messageVo.getCurrentRetryCount(), messageVo);
} catch (Exception e) {
log.error("将异常消息存储到mongodb失败,消息数据:" + messageVo.toString(), e);
}
}
}
- RetryController.java
import cn.hutool.core.util.StrUtil;
import cn.hutool.json.JSONObject;
import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.core.MessageProperties;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RequestParam;
import org.springframework.web.bind.annotation.RestController;
import run.siyuan.common.rabbitmq.Constants;
@Slf4j
@RestController
@RequestMapping("/retry")
public class RetryController {
@Autowired
private RabbitTemplate rabbitTemplate;
@GetMapping(value = "/service/message")
public String consumerFailQueue(@RequestParam(required = false, defaultValue = "1") Integer id) {
JSONObject json = new JSONObject();
json.set("id", id);
json.set("name", "消息名称");
json.set("time", System.currentTimeMillis());
String msg = StrUtil.format("消息发送时间:{} 消息数据:{}", System.currentTimeMillis(), json);
log.info(msg);
rabbitTemplate.convertAndSend(Constants.SERVICE_QUEUE, json);
log.info("消息发送完成时间:{}", System.currentTimeMillis());
return "success";
}
}



