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

如何为Spring Boot ReactiveRedis实现自定义RedisListener订阅注解

基于Spring Boot ReactiveRedis实现自定义@RedisListener注解的发布订阅机制

现状说明

目前已通过ReactiveRedisTemplate实现基础的Redis发布订阅功能,代码运行正常,能在订阅回调中获取数据:

@Autowired
private ReactiveRedisOperations<String, Object> reactiveRedisTemplate;
 
@PostConstruct
public void init(){
    System.out.println("*****SampleLoader***** initialized - SampleLoader");

    this.reactiveRedisTemplate.listenTo(ChannelTopic.of("some-topic")).subscribe(data -> {
            System.out.println(
                    "data.getChannel():-" + data.getChannel() + ":" + "data.getMessage():-" + data.getMessage());
    });
}

需求目标

实现自定义@RedisListener注解,使得标注该注解的方法能够自动订阅指定Redis主题,并在收到消息时自动触发方法回调,接收通道和消息数据。期望的注解使用方式如下:

@Target({ ElementType.TYPE, ElementType.METHOD, ElementType.ANNOTATION_TYPE })
@Retention(RetentionPolicy.RUNTIME)
@Documented
@Repeatable(RedisListeners.class)
public @interface RedisListener{
}

@Target({ ElementType.TYPE, ElementType.METHOD, ElementType.ANNOTATION_TYPE })
@Retention(RetentionPolicy.RUNTIME)
@Documented
public @interface RedisListeners {
    RedisListener[] value();
}

// 使用示例
@RedisListener("some-topic")
public void redisData(String channel, Object data){
    // 处理消息逻辑
}

实现步骤与代码

1. 完善自定义注解定义

首先为@RedisListener添加主题属性,支持单个或多个主题配置:

import java.lang.annotation.*;

@Target({ElementType.METHOD})
@Retention(RetentionPolicy.RUNTIME)
@Documented
@Repeatable(RedisListeners.class)
public @interface RedisListener {
    // 配置订阅的Redis主题,支持多个主题
    String[] value() default {};
}

@Target({ElementType.METHOD})
@Retention(RetentionPolicy.RUNTIME)
@Documented
public @interface RedisListeners {
    RedisListener[] value();
}

2. 实现注解处理器

编写一个Spring组件,在容器初始化完成后扫描所有Bean中的方法,识别标注了@RedisListener的方法,并自动创建Redis订阅逻辑:

import org.springframework.context.ApplicationContext;
import org.springframework.context.ApplicationListener;
import org.springframework.context.event.ContextRefreshedEvent;
import org.springframework.data.redis.core.ReactiveRedisOperations;
import org.springframework.data.redis.listener.ChannelTopic;
import org.springframework.stereotype.Component;
import java.lang.reflect.Method;

@Component
public class RedisAnnotationListenerProcessor implements ApplicationListener<ContextRefreshedEvent> {

    private final ReactiveRedisOperations<String, Object> reactiveRedisTemplate;

    // 构造注入ReactiveRedisTemplate
    public RedisAnnotationListenerProcessor(ReactiveRedisOperations<String, Object> reactiveRedisTemplate) {
        this.reactiveRedisTemplate = reactiveRedisTemplate;
    }

    @Override
    public void onApplicationEvent(ContextRefreshedEvent event) {
        ApplicationContext context = event.getApplicationContext();
        // 遍历容器中所有Bean
        for (String beanName : context.getBeanDefinitionNames()) {
            Object bean = context.getBean(beanName);
            // 遍历Bean的所有方法
            for (Method method : bean.getClass().getDeclaredMethods()) {
                // 处理单个@RedisListener注解
                RedisListener[] singleListeners = method.getAnnotationsByType(RedisListener.class);
                for (RedisListener listener : singleListeners) {
                    bindListenerToTopics(listener.value(), bean, method);
                }
                // 处理@RedisListeners批量注解
                RedisListeners batchListeners = method.getAnnotation(RedisListeners.class);
                if (batchListeners != null) {
                    for (RedisListener listener : batchListeners.value()) {
                        bindListenerToTopics(listener.value(), bean, method);
                    }
                }
            }
        }
    }

    /**
     * 为指定主题绑定方法回调
     */
    private void bindListenerToTopics(String[] topics, Object bean, Method method) {
        for (String topic : topics) {
            reactiveRedisTemplate.listenTo(ChannelTopic.of(topic))
                    .subscribe(message -> {
                        try {
                            // 调用标注的方法,传递通道名和消息数据
                            method.invoke(bean, message.getChannel(), message.getMessage());
                        } catch (Exception e) {
                            // 可根据业务需求替换为日志记录或异常处理逻辑
                            e.printStackTrace();
                        }
                    });
        }
    }
}

3. 使用示例

在Spring组件的方法上标注@RedisListener即可自动订阅指定主题:

import org.springframework.stereotype.Component;

@Component
public class RedisMessageHandler {

    // 订阅单个主题
    @RedisListener("some-topic")
    public void handleSingleTopic(String channel, Object data) {
        System.out.println("收到单主题消息:通道=" + channel + ",数据=" + data);
    }

    // 同时订阅多个主题
    @RedisListener({"topic-user", "topic-order"})
    public void handleMultipleTopics(String channel, Object data) {
        System.out.println("收到多主题消息:通道=" + channel + ",数据=" + data);
    }

    // 使用@RedisListeners重复标注不同主题
    @RedisListeners({
            @RedisListener("topic-notify"),
            @RedisListener("topic-log")
    })
    public void handleBatchTopics(String channel, Object data) {
        System.out.println("收到批量主题消息:通道=" + channel + ",数据=" + data);
    }
}

注意事项

  • 确保ReactiveRedisTemplate已正确配置,消息序列化方式与发布端一致,避免数据解析异常
  • 标注@RedisListener的方法参数需与处理器中传递的参数类型匹配(String channel, Object data),若需自定义参数类型,需在处理器中添加类型转换逻辑
  • 订阅逻辑基于Reactive异步模型,需注意方法执行的线程安全问题
  • 异常处理可根据业务需求优化,比如集成日志框架记录错误信息

内容的提问来源于stack exchange,提问作者Ganesh Gudghe

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 18:32:28