yudaocode/ruoyi-vue-pro:Redis消息队列实战
·
yudaocode/ruoyi-vue-pro:Redis消息队列实战
引言
在现代分布式系统中,消息队列(Message Queue)是实现异步通信、解耦服务、削峰填谷的关键技术。Redis作为高性能的内存数据库,不仅提供缓存功能,其List数据结构更是实现轻量级消息队列的绝佳选择。本文将深入探讨如何在芋道RuoYi-Vue-Pro项目中实战Redis消息队列,解决实际业务场景中的异步处理需求。
Redis消息队列核心原理
数据结构选择
Redis提供了多种数据结构来实现消息队列,最常用的是List结构:
核心命令解析
| 命令 | 作用 | 使用场景 |
|---|---|---|
LPUSH |
从左侧插入消息 | 生产者发送消息 |
RPOP |
从右侧取出消息 | 消费者非阻塞获取 |
BRPOP |
阻塞式右侧取出 | 消费者实时监听 |
LLEN |
获取队列长度 | 监控队列状态 |
LTRIM |
修剪队列 | 维护队列大小 |
项目集成实战
1. 环境配置
首先确保Redis依赖已正确配置:
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-data-redis</artifactId>
</dependency>
2. Redis配置类
@Configuration
public class RedisMessageQueueConfig {
@Bean
public RedisTemplate<String, Object> redisTemplate(RedisConnectionFactory factory) {
RedisTemplate<String, Object> template = new RedisTemplate<>();
template.setConnectionFactory(factory);
template.setKeySerializer(new StringRedisSerializer());
template.setValueSerializer(new GenericJackson2JsonRedisSerializer());
return template;
}
}
3. 消息队列服务实现
@Service
@Slf4j
public class RedisMessageQueueService {
@Autowired
private RedisTemplate<String, Object> redisTemplate;
private static final String QUEUE_PREFIX = "message:queue:";
/**
* 发送消息到队列
*/
public void sendMessage(String queueName, Object message) {
String key = QUEUE_PREFIX + queueName;
redisTemplate.opsForList().leftPush(key, message);
log.info("消息已发送到队列 {}: {}", queueName, message);
}
/**
* 阻塞式接收消息
*/
public Object receiveMessage(String queueName, long timeout) {
String key = QUEUE_PREFIX + queueName;
return redisTemplate.opsForList().rightPop(key, timeout, TimeUnit.SECONDS);
}
/**
* 获取队列长度
*/
public long getQueueSize(String queueName) {
String key = QUEUE_PREFIX + queueName;
Long size = redisTemplate.opsForList().size(key);
return size != null ? size : 0;
}
}
4. 生产者示例
@RestController
@RequestMapping("/message")
public class MessageProducerController {
@Autowired
private RedisMessageQueueService messageQueueService;
@PostMapping("/send")
public ApiResult<String> sendMessage(@RequestBody MessageDTO message) {
// 业务验证
if (StringUtils.isEmpty(message.getContent())) {
return ApiResult.error("消息内容不能为空");
}
// 发送到消息队列
messageQueueService.sendMessage("email_queue", message);
return ApiResult.success("消息发送成功");
}
}
5. 消费者示例
@Component
@Slf4j
public class EmailMessageConsumer {
@Autowired
private RedisMessageQueueService messageQueueService;
@Autowired
private EmailService emailService;
@PostConstruct
public void startConsumer() {
new Thread(this::consumeMessages).start();
}
private void consumeMessages() {
while (true) {
try {
Object message = messageQueueService.receiveMessage("email_queue", 30);
if (message != null) {
processMessage((MessageDTO) message);
}
} catch (Exception e) {
log.error("消息处理异常", e);
}
}
}
private void processMessage(MessageDTO message) {
try {
// 发送邮件逻辑
emailService.sendEmail(message.getTo(), message.getSubject(), message.getContent());
log.info("邮件发送成功: {}", message);
} catch (Exception e) {
log.error("邮件发送失败: {}", message, e);
// 重试机制或死信队列处理
}
}
}
高级特性实现
1. 延迟队列实现
public class DelayedMessageQueue {
@Autowired
private RedisTemplate<String, Object> redisTemplate;
private static final String DELAYED_QUEUE = "delayed:queue";
private static final String READY_QUEUE_PREFIX = "ready:queue:";
/**
* 添加延迟消息
*/
public void addDelayedMessage(String queueName, Object message, long delaySeconds) {
double score = System.currentTimeMillis() + delaySeconds * 1000;
redisTemplate.opsForZSet().add(DELAYED_QUEUE,
new DelayedMessage(queueName, message), score);
}
/**
* 检查并转移到期消息
*/
@Scheduled(fixedRate = 1000)
public void checkDelayedMessages() {
Set<ZSetOperations.TypedTuple<Object>> messages =
redisTemplate.opsForZSet().rangeByScoreWithScores(
DELAYED_QUEUE, 0, System.currentTimeMillis());
for (ZSetOperations.TypedTuple<Object> tuple : messages) {
DelayedMessage delayedMessage = (DelayedMessage) tuple.getValue();
redisTemplate.opsForList().leftPush(
READY_QUEUE_PREFIX + delayedMessage.getQueueName(),
delayedMessage.getMessage());
redisTemplate.opsForZSet().remove(DELAYED_QUEUE, delayedMessage);
}
}
}
2. 消息确认机制
public class AckMessageQueueService {
private static final String PROCESSING_QUEUE = "processing:queue:";
private static final long PROCESSING_TIMEOUT = 300; // 5分钟
/**
* 带确认的消息消费
*/
public Object consumeWithAck(String queueName) {
Object message = messageQueueService.receiveMessage(queueName, 30);
if (message != null) {
// 将消息放入处理中队列
String messageId = generateMessageId(message);
redisTemplate.opsForValue().set(
PROCESSING_QUEUE + messageId,
message, PROCESSING_TIMEOUT, TimeUnit.SECONDS);
}
return message;
}
/**
* 消息确认
*/
public void ackMessage(String messageId) {
redisTemplate.delete(PROCESSING_QUEUE + messageId);
}
/**
* 消息重投递
*/
@Scheduled(fixedRate = 60000)
public void redeliverTimeoutMessages() {
Set<String> keys = redisTemplate.keys(PROCESSING_QUEUE + "*");
for (String key : keys) {
Long ttl = redisTemplate.getExpire(key);
if (ttl != null && ttl < 60) {
Object message = redisTemplate.opsForValue().get(key);
String queueName = extractQueueNameFromKey(key);
messageQueueService.sendMessage(queueName, message);
redisTemplate.delete(key);
}
}
}
}
性能优化策略
1. 批量处理优化
public class BatchMessageProcessor {
/**
* 批量发送消息
*/
public void batchSend(String queueName, List<Object> messages) {
String key = QUEUE_PREFIX + queueName;
redisTemplate.executePipelined((RedisCallback<Object>) connection -> {
for (Object message : messages) {
connection.lPush(serializeKey(key), serializeValue(message));
}
return null;
});
}
/**
* 批量消费消息
*/
public List<Object> batchReceive(String queueName, int batchSize) {
String key = QUEUE_PREFIX + queueName;
List<Object> messages = new ArrayList<>();
for (int i = 0; i < batchSize; i++) {
Object message = messageQueueService.receiveMessage(queueName, 1);
if (message == null) break;
messages.add(message);
}
return messages;
}
}
2. 内存优化配置
spring:
redis:
lettuce:
pool:
max-active: 8
max-idle: 8
min-idle: 0
max-wait: -1ms
timeout: 2000ms
监控与告警
1. 队列监控
@Component
public class QueueMonitor {
@Autowired
private RedisMessageQueueService messageQueueService;
@Scheduled(fixedRate = 30000)
public void monitorQueues() {
Map<String, Long> queueStats = new HashMap<>();
// 监控所有业务队列
String[] queues = {"email_queue", "sms_queue", "push_queue", "order_queue"};
for (String queue : queues) {
long size = messageQueueService.getQueueSize(queue);
queueStats.put(queue, size);
// 告警逻辑
if (size > 1000) {
sendAlert("队列 " + queue + " 积压严重,当前数量: " + size);
}
}
log.info("队列监控统计: {}", queueStats);
}
}
2. 性能指标采集
@Aspect
@Component
@Slf4j
public class QueuePerformanceAspect {
@Around("execution(* com.example.service..*.*(..))")
public Object monitorPerformance(ProceedingJoinPoint joinPoint) throws Throwable {
long startTime = System.currentTimeMillis();
Object result = joinPoint.proceed();
long endTime = System.currentTimeMillis();
String methodName = joinPoint.getSignature().getName();
long duration = endTime - startTime;
// 记录到监控系统
Metrics.recordTiming("queue_processing_time", duration,
"method", methodName);
return result;
}
}
最佳实践总结
1. 设计原则
- 单一职责: 每个队列只处理一种类型的消息
- 适度超时: 设置合理的消息处理超时时间
- 重试机制: 实现消息重试和死信队列
- 监控告警: 实时监控队列状态和性能指标
2. 适用场景
| 场景 | 推荐方案 | 注意事项 |
|---|---|---|
| 邮件发送 | Redis List + 批量处理 | 注意邮件服务商的频率限制 |
| 短信通知 | 延迟队列 + 流量控制 | 避免触发运营商风控 |
| 订单处理 | 事务消息 + 幂等性 | 保证消息处理的可靠性 |
| 日志收集 | 批量处理 + 压缩 | 优化网络传输效率 |
3. 避坑指南
- 内存管理: 定期清理已完成的消息,避免内存泄漏
- 网络超时: 配置合理的Redis连接超时时间
- 序列化: 使用JSON序列化保证消息的可读性和兼容性
- 集群部署: 在Redis集群环境下注意key的分布策略
通过本文的实战指南,您可以在芋道RuoYi-Vue-Pro项目中快速搭建稳定高效的Redis消息队列系统,为业务异步化处理提供强有力的技术支撑。
更多推荐


所有评论(0)