yudaocode/ruoyi-vue-pro:Redis消息队列实战

【免费下载链接】ruoyi-vue-pro 🔥 官方推荐 🔥 RuoYi-Vue 全新 Pro 版本,优化重构所有功能。基于 Spring Boot + MyBatis Plus + Vue & Element 实现的后台管理系统 + 微信小程序,支持 RBAC 动态权限、数据权限、SaaS 多租户、Flowable 工作流、三方登录、支付、短信、商城、CRM、ERP、AI 等功能。你的 ⭐️ Star ⭐️,是作者生发的动力! 【免费下载链接】ruoyi-vue-pro 项目地址: https://gitcode.com/yudaocode/ruoyi-vue-pro

引言

在现代分布式系统中,消息队列(Message Queue)是实现异步通信、解耦服务、削峰填谷的关键技术。Redis作为高性能的内存数据库,不仅提供缓存功能,其List数据结构更是实现轻量级消息队列的绝佳选择。本文将深入探讨如何在芋道RuoYi-Vue-Pro项目中实战Redis消息队列,解决实际业务场景中的异步处理需求。

Redis消息队列核心原理

数据结构选择

Redis提供了多种数据结构来实现消息队列,最常用的是List结构:

mermaid

核心命令解析

命令 作用 使用场景
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. 避坑指南

  1. 内存管理: 定期清理已完成的消息,避免内存泄漏
  2. 网络超时: 配置合理的Redis连接超时时间
  3. 序列化: 使用JSON序列化保证消息的可读性和兼容性
  4. 集群部署: 在Redis集群环境下注意key的分布策略

通过本文的实战指南,您可以在芋道RuoYi-Vue-Pro项目中快速搭建稳定高效的Redis消息队列系统,为业务异步化处理提供强有力的技术支撑。

【免费下载链接】ruoyi-vue-pro 🔥 官方推荐 🔥 RuoYi-Vue 全新 Pro 版本,优化重构所有功能。基于 Spring Boot + MyBatis Plus + Vue & Element 实现的后台管理系统 + 微信小程序,支持 RBAC 动态权限、数据权限、SaaS 多租户、Flowable 工作流、三方登录、支付、短信、商城、CRM、ERP、AI 等功能。你的 ⭐️ Star ⭐️,是作者生发的动力! 【免费下载链接】ruoyi-vue-pro 项目地址: https://gitcode.com/yudaocode/ruoyi-vue-pro

Logo

Agent 垂直技术社区,欢迎活跃、内容共建。

更多推荐