You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何将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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.28 09:43:16