如何将RabbitMQ消息存入Redis缓存/数据库?求Demo代码实现
微服务场景下RabbitMQ消息监听+Redis存储Demo
一、基础依赖(Spring Boot)
在pom.xml中引入所需依赖:
<dependencies> <!-- Spring Boot RabbitMQ 依赖 --> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-amqp</artifactId> </dependency> <!-- Spring Boot Redis 依赖 --> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-data-redis</artifactId> </dependency> <!-- Jackson 用于JSON序列化 --> <dependency> <groupId>com.fasterxml.jackson.core</groupId> <artifactId>jackson-databind</artifactId> </dependency> </dependencies>
二、监听服务实现
1. 配置文件(application.yml)
spring: rabbitmq: host: localhost port: 5672 username: guest password: guest virtual-host: / redis: host: localhost port: 6379 password: "" database: 0
2. RabbitMQ队列配置
import org.springframework.amqp.core.Queue; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; @Configuration public class RabbitMqConfig { // 定义监听的队列名称 public static final String SCHEDULER_UPDATE_QUEUE = "scheduler.update.queue"; @Bean public Queue schedulerUpdateQueue() { // 持久化队列 return new Queue(SCHEDULER_UPDATE_QUEUE, true); } }
3. 消息消费者与Redis存储逻辑
import org.springframework.amqp.rabbit.annotation.RabbitListener; import org.springframework.data.redis.core.StringRedisTemplate; import org.springframework.stereotype.Component; import java.time.LocalDateTime; import java.time.format.DateTimeFormatter; @Component public class SchedulerUpdateListener { private final StringRedisTemplate stringRedisTemplate; // Redis键前缀,用于区分不同调度服务的记录 private static final String REDIS_KEY_PREFIX = "scheduler:last_update:"; // 时间格式化器 private static final DateTimeFormatter DATE_TIME_FORMATTER = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss"); public SchedulerUpdateListener(StringRedisTemplate stringRedisTemplate) { this.stringRedisTemplate = stringRedisTemplate; } @RabbitListener(queues = RabbitMqConfig.SCHEDULER_UPDATE_QUEUE) public void handleSchedulerUpdate(String schedulerName) { // 获取当前时间作为最后更新时间 String lastUpdateTime = LocalDateTime.now().format(DATE_TIME_FORMATTER); // 将调度服务名称和最后更新时间存入Redis,键为前缀+服务名,值为更新时间 stringRedisTemplate.opsForValue().set(REDIS_KEY_PREFIX + schedulerName, lastUpdateTime); // 可选:打印日志确认操作 System.out.printf("已记录调度服务[%s]的最后更新时间:%s%n", schedulerName, lastUpdateTime); } }
三、调度服务实现(两个示例)
调度服务1:SchedulerServiceA
import org.springframework.amqp.rabbit.core.RabbitTemplate; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Component; @Component public class SchedulerServiceA { private final RabbitTemplate rabbitTemplate; // 当前服务名称 private static final String SERVICE_NAME = "scheduler-service-a"; public SchedulerServiceA(RabbitTemplate rabbitTemplate) { this.rabbitTemplate = rabbitTemplate; } // 模拟每30秒执行一次调度任务,发送更新消息 @Scheduled(fixedRate = 30000) public void executeScheduleTask() { // 发送服务名称到RabbitMQ队列 rabbitTemplate.convertAndSend(RabbitMqConfig.SCHEDULER_UPDATE_QUEUE, SERVICE_NAME); System.out.printf("调度服务[%s]已发送更新消息%n", SERVICE_NAME); } }
调度服务2:SchedulerServiceB
import org.springframework.amqp.rabbit.core.RabbitTemplate; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Component; @Component public class SchedulerServiceB { private final RabbitTemplate rabbitTemplate; // 当前服务名称 private static final String SERVICE_NAME = "scheduler-service-b"; public SchedulerServiceB(RabbitTemplate rabbitTemplate) { this.rabbitTemplate = rabbitTemplate; } // 模拟每45秒执行一次调度任务,发送更新消息 @Scheduled(fixedRate = 45000) public void executeScheduleTask() { // 发送服务名称到RabbitMQ队列 rabbitTemplate.convertAndSend(RabbitMqConfig.SCHEDULER_UPDATE_QUEUE, SERVICE_NAME); System.out.printf("调度服务[%s]已发送更新消息%n", SERVICE_NAME); } }
关键启动配置
所有服务的启动类需要添加@EnableRabbit和@EnableScheduling注解,开启RabbitMQ监听和定时任务功能:
import org.springframework.boot.SpringApplication; import org.springframework.boot.autoconfigure.SpringBootApplication; import org.springframework.scheduling.annotation.EnableScheduling; @SpringBootApplication @EnableRabbit @EnableScheduling public class Application { public static void main(String[] args) { SpringApplication.run(Application.class, args); } }
注意事项
- 确保RabbitMQ和Redis服务已启动并配置正确。
- 若需传递更复杂的消息内容,可定义DTO类,通过RabbitMQ的JSON序列化器实现对象的发送与接收。
内容的提问来源于stack exchange,提问作者Nifi Logstash
相关产品推荐
相关产品推荐

