1. 聊清楚业务场景:Redis消息订阅到底在解决什么问题
如果你维护过一个多实例部署的服务,或者一整套微服务集群,大概率遇到过这种尴尬:订单服务更新了一条商品数据,但另外几个服务实例的本地缓存里还是旧数据,用户下单时看到了严重不一致的价格。又或者,某个新节点想要上线,希望集群里其他节点立刻感知到它的存在,而不是傻等30秒一查的注册中心心跳。
这种"一个节点发生了事情,要立刻让其他所有节点都知道"的需求,最朴素的做法就是轮询:每个实例定时去数据库刷一遍,比对是否有变化。轮询实现简单,但代价很高——数据量大了以后,每一次轮询都是对数据库的无效压力,而且时效性永远受制于轮询间隔。另一个做法是让通知发起方主动调用其他节点的HTTP接口,但你要维护一份"所有节点地址列表",节点扩缩容时还得动态更新,稍不注意就是一个定时炸弹。
Redis作为消息订阅通道,恰好能在这个层面试着解决一类问题。它天然支持发布订阅模式,一个客户端向某个频道发布消息,所有订阅了该频道的客户端都会立刻收到。Redis在绝大多数Java后端项目里本来就已经是基础设施,不需要额外引入新的中间件,不需要单独运维一套集群,Spring Data Redis又已经把订阅相关的API封装得比较完整。这也是我为什么在实际项目中遇到"轻量广播"类需求时,会优先考虑它而不是一上来就引入RabbitMQ或Kafka。
这篇文章我会从场景分析、环境准备、核心原理、可落地的Demo、进阶玩法和实际踩坑几个方面,把"用Spring订阅Redis消息"这件事完整讲透。适合那些已经会基本Spring Boot开发、想把Redis消息订阅功能用起来的后端开发者,也适合正在做技术选型、想搞明白Redis Pub/Sub边界的人。
1.1 为什么会有这个需求:从轮询和HTTP通知的尴尬说起
先假设一个最简单的多实例场景:三个Java服务实例连同一个数据库,每个实例内部都有一层Caffeine本地缓存,用来缓存用户信息。用户修改昵称后,数据库被更新了,但是另外两个实例的Caffeine里还是旧昵称。如果缓存TTL设置得比较长,用户看到旧昵称的时间就会很长;如果TTL设置得短,本地缓存又起不到应有的加速效果。
这种"本地缓存失效"问题,天然需要一个广播通道:谁修改了数据,就通知所有实例把对应key的缓存删掉。HTTP互相调用做不到完全可靠,还得解决服务发现;数据库轮询又太浪费;Redis是现有基础设施,Pub/Sub一条命令就能广播,成本最低。类似的场景还包括:配置中心推送配置变更、停机维护前通知所有节点拒绝流量、在线用户的状态变更通知等等。
我后来做项目时总结了一条经验:凡是"集群内所有实例都需要知道同一件事"且"丢了这条消息不至于立刻出大事故"的通知,都可以用Redis Pub/Sub来承担。它解决的是广播通知的问题,不是可靠投递的问题,这条边界一定要先刻在脑子里。
1.2 Pub/Sub的核心机制:频道、发布者、订阅者
Redis的发布订阅模型非常简单,只有三个角色:发布者、频道、订阅者。发布者通过PUBLISH命令向某个频道发消息,订阅者通过SUBSCRIBE命令订阅某个频道。一条消息发出后,Redis会实时转发给所有订阅者,没有订阅者的消息就直接被丢弃,绝不会为"未来可能的订阅者"缓存任何东西。
与之配套的还有PSUBSCRIBE命令,它支持模式匹配订阅,比如订阅"cache:*"就能接收到所有以"cache:"开头的频道消息。这个功能在需要按业务域批量处理消息时非常有用,后面我会专门展开。退订则对应UNSUBSCRIBE和PUNSUBSCRIBE。
这里有个特别容易误解的点:Redis Pub/Sub不使用索引,不同于传统的键值读写,它不保存数据。消息一旦被转发完成,在Redis内存里不会留下任何痕迹。这一点和Redis Streams完全不同,而且恰恰决定了它的适用边界——只适合实时在线通知,不适合"你下线了,等你回来我给你补发"这种场景。
1.3 哪些场景适合用Pub/Sub:和MQ的边界划分
很多人在选型时容易走极端,要么所有通知都上MQ,要么觉得Redis Pub/Sub能通吃。我的经验是:这是一个分层决策的问题,做之前先画个表对照一下。
| 特性 | Redis Pub/Sub | Redis Streams(5.0+) | RabbitMQ / Kafka |
|---|---|---|---|
| 消息持久化 | 无,发完即走 | 有,可回溯 | 有 |
| 消费确认 | 无,不重试 | 有,ACK机制 | 有 |
| 消息堆积 | 不支持,在线实时推送 | 支持按范围读取 | 支持 |
| 消费者分组 | 无,天然广播 | 支持消费组竞争 | 支持 |
| 集成复杂度 | 低,Spring封装完善 | 中等 | 高,需要独立服务 |
| 典型场景 | 缓存失效、状态广播 | 任务队列、可靠日志流 | 核心业务异步解耦 |
所以你问我Redis Pub/Sub适合什么人用,我会说:适合那些"消息丢了也能接受、但希望实时性特别好"的场景。缓存失效广播是最典型的,就算某条失效通知丢了,充其量是多等一个TTL,数据不会有错。在线状态通知也能接受,用户下线了通知丢了,下次心跳一报照样状态准确。反过来,如果是订单创建后的积压处理、支付结果通知给客户端、跨系统数据同步,这些场景对可靠性和确认机制有硬性要求,不应该拿Pub/Sub硬扛。这些判断不是Redis本身的缺陷,而是机制特性决定的,选型时尊重它就好。
2. 环境准备:Redis服务端与Spring Boot项目的最小配置
2.1 一分钟把Redis跑起来:本地与Docker
先说服务端。本地开发最常见的方式是直接用Docker起一个,干净利落不污染宿主机:
docker run -d --name redis -p 6379:6379 redis:7-alpine这条命令用官方Redis 7镜像启动一个容器,并映射默认端口6379。如果你在Windows上开发又不想用Docker,直接去Redis官网下载Windows压缩包,解压后双击redis-server.exe也能运行,注意Windows版的版本一般滞后于Linux官方版,但做开发足够了。
如果你的现有环境是Redis主从架构,比如用Docker部署了一主一从,那么在搭建Pub/Sub功能时有一条重要经验:订阅连接尽量指向主节点。原因在于常规主从复制模式下,PUBLISH消息在主从之间的复制语义比较特殊,订阅从库容易遇到收不到消息或者连接切换后订阅失效的情况。后面讲连接恢复坑的时候我会回来详谈。
调试阶段我建议装一个可视化客户端,RedisInsight或者Another Redis Desktop Manager都行。虽然Pub/Sub本身不存数据,在这些工具里看不到消息历史,但是你可以直接手动发一条测试消息,比用命令行方便不少。Redis的命令行工具redis-cli也够用,redis-cli publish cache:invalidate "test"一敲,测试消息就发出去了。
2.2 引入依赖与连接配置:版本差异要注意
Spring Boot项目里引入Redis的依赖只有一个:
<dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-data-redis</artifactId> </dependency>这个starter会传递引入Spring Data Redis和对应的连接器。默认情况下Spring Boot 2.x和3.x都使用Lettuce作为Redis客户端,Lettuce基于Netty实现,连接管理和异步支持都做得不错,一般不需要换Jedis,除非有特殊的历史包袱。
再来看配置。这里有一个版本差异坑:Spring Boot 2.x的连接配置前缀是spring.redis,到了Spring Boot 3.x改成了spring.data.redis。如果你从2.x升级到3.x,忘了改前缀会导致配置静默失效,本地能连上(因为默认localhost),一上生产就连接不上。我的示例以Spring Boot 3.x为主:
spring: data: redis: host: 127.0.0.1 port: 6379 password: database: 0 lettuce: pool: enabled: true max-active: 16 max-idle: 8 min-idle: 2连接池默认是不开启的,如果你只是做简单的消息收发,不开启也能跑。但考虑到RedisTemplate做读写、Pub/Sub订阅连接、还有多线程场景下同时获取连接,建议显式打开连接池,避免在高并发下出现"无法获取连接"的报错。
2.3 接收端的数据通道:StringRedisTemplate的作用
Spring Data Redis对RedisTemplate和StringRedisTemplate做了区分:RedisTemplate默认使用JdkSerializationRedisSerializer,会把对象序列化成JDK二进制字节后再写入Redis;StringRedisTemplate则强制key和value都使用StringRedisSerializer,读写数据时直接按UTF-8字节处理。
做消息订阅时,我强烈建议用StringRedisTemplate作为发送通道。原因后面讲序列化坑的时候会详细说,这里先记住结论:Pub/Sub的订阅者和发布者之间传递的本质上就是字节流,用StringRedisTemplate可以让这个字节流保持最朴素的UTF-8文本形态,接收端拿到后直接转字符串就完事,完全绕开JDK序列化带来的兼容性问题。
在Spring Boot项目中,StringRedisTemplate是自动配置好的,直接注入就能用,甚至不需要额外注册Bean。很多初学者上来就是RedisTemplate<String, String>再配一堆Serializer,结果发送端和接收端序列化策略不一致,消息一发出去就变成一堆看不懂的乱码。用StringRedisTemplate是最省心的入口。
3. 从原生订阅到Spring容器:理解RedisMessageListenerContainer
3.1 原生Jedis订阅为什么麻烦:阻塞式API的先天问题
如果你不借助Spring的高层封装,直接用Jedis或Lettuce写订阅代码,会发现一个很别扭的地方:Redis的订阅命令是阻塞式的。client.subscribe(listener, channel)一旦执行,当前线程就会一直阻塞在这里等待消息,线程内后面任何代码都执行不到。所以你必须为每一个订阅单独开一条线程,底层连接也必须单独占用,不能和普通的读写操作共用一个连接。
写出来大概是这样:
Jedis jedis = new Jedis("127.0.0.1", 6379); new Thread(() -> { jedis.subscribe(new JedisPubSub() { @Override public void onMessage(String channel, String message) { System.out.println("收到消息: " + message); } }, "cache:invalidate"); }).start();这段代码能跑,但你很快就会遇到一连串问题:线程怎么管理?连接断开后怎么重连?多个监听频道时要不要每个都开线程?应用关闭时这些阻塞线程怎么优雅停掉?如果这些事都让业务代码来处理,代码很快就会失控。这个道理很像ServerSocket和Servlet容器的区别——裸Socket能做,但让容器帮你管理生命周期才是正路。
3.2 Spring容器帮你管了什么:连接、线程与生命周期
RedisMessageListenerContainer就是Spring Data Redis提供的那个"Servlet容器"。它是一个Spring Bean,实现了SmartLifecycle接口,也就是说它在IoC容器启动后会自行启动,应用关闭时会自行停止。它内部帮你处理了好几件琐碎的事:
- 每个订阅任务独占一条Redis连接,连接由它统一从连接工厂获取,不会和普通RedisTemplate操作抢连接;
- 每个监听器对应一个阻塞订阅线程,这个线程由TaskExecutor执行器统一调度;
- 连接断开后会自动重连并重新提交订阅,不需要你手动打补丁;
- 应用关闭时能优雅地释放连接、停止线程,避免资源泄漏。
它的核心API也不复杂:addMessageListener(MessageListener, Topic)负责注册监听器和它关心的频道,removeMessageListener负责移除。底层的MessageListener接口只有一个方法:
void onMessage(Message message, byte[] pattern);Message对象里有三个关键信息:消息体(getBody())、频道(getChannel())、还有消息在模式订阅时命中的pattern(getPattern())。Spring把原始字节交给你的Listener,你拥有完全的控制权,而不会像某些框架那样自作主张帮你反序列化然后把异常藏起来。
3.3 一条消息从PUBLISH到onMessage的完整路径
把端到端的链路画出来,能帮助你以后排查问题。我描述一下这条链路:假设某个服务实例执行了stringRedisTemplate.convertAndSend("cache:invalidate", "book:1"),数据先被序列化成UTF-8字节,通过Lettuce的连接发送到Redis服务器。Redis服务器根据频道名找到所有订阅了该频道的连接,把消息体、原始频道字节一起推送到这些连接。
此时,每个应用实例里的RedisMessageListenerContainer,其内部为这个订阅分配的那个阻塞线程会收到数据,然后Dispatcher线程池将消息派发给对应的MessageListener。如果监听的是PatternTopic,pattern字段会带上匹配用的模式字节;如果监听的是普通ChannelTopic,pattern字段就是空值。
所以日后排查"消息没收到"的时候,审视的环节就是:发布端有没有真正发到Redis?Redis服务器有没有这个频道对应的订阅连接?接收端的订阅连接有没有保持在线?容器Dispatcher有没有把消息正确派发?一条链路拆开,逐个环节验证,问题就藏不住。
4. 写一个能落地的Demo:基于Pub/Sub的缓存失效广播
4.1 业务设定:本地缓存如何被一条消息统一失效
这里我拿一个实际做过的商业项目来举例。我们的几个微服务实例都用了Caffeine做本地缓存,热点商品信息缓存5分钟。运营修改商品价格后,数据库立刻更新,但其他实例的Caffeine里还存着旧价格,最长可能5分钟不更新。客服后台那边用户能直观看到新旧价格不一致,体验很差。
最终方案就是引入Redis Pub/Sub做缓存失效广播。修改商品时先更新数据库,再向cache:invalidate频道发一条消息,消息的内容是失效的key,比如"product:1001"。所有监听该频道的实例收到后,调用本地Caffeine的invalidate(key)删掉这一条。这个设计一上线,缓存不一致的窗口就从"最多5分钟"缩小到了"毫秒级"。
这个方案里消息丢失是可以接受的:就算某一条失效消息丢了,最坏情况也就是这个实例的缓存继续用到自然过期,不会产生数据错乱。这正是我之前说的边界——适合、且只适合这类对可靠性要求不高的广播。
4.2 监听器实现:直接处理字节,绕过序列化魔咒
监听器最简单也最稳妥的写法,是直接实现MessageListener接口,拿到字节后手动转String。不依赖任何容器层面的消息转换器:
@Component public class CacheInvalidateListener implements MessageListener { private static final Logger log = LoggerFactory.getLogger(CacheInvalidateListener.class); private final CacheManager cacheManager; public CacheInvalidateListener(CacheManager cacheManager) { this.cacheManager = cacheManager; } @Override public void onMessage(Message message, byte[] pattern) { String channel = new String(message.getChannel(), StandardCharsets.UTF_8); String key = new String(message.getBody(), StandardCharsets.UTF_8); log.info("收到缓存失效广播,channel={}, key={}", channel, key); cacheManager.invalidate(key); } }这段代码看起来平平无奇,但有一个很重要的设计选择:为什么不用Spring提供的MessageListenerAdapter直接让方法接收一个String参数?原因在于MessageListenerAdapter的底层依赖消息转换器,转换器的序列化策略一旦和发送端不一致,轻则收到乱码,重则直接抛SerializationException。直接处理字节是最稳定的方式,把new String这一步自己控制住,序列化这个变量就从系统里排除了。
当然,如果你们的统一规范里已经约定好所有Redis操作都使用JSON序列化,那么搞一个全局的RedisMessageListenerContainer配置注入GenericJackson2JsonRedisSerializer也是可以的。但我不建议在项目初期就这么干,先跑通消息链路,再考虑优雅不优雅。
4.3 容器配置:Topic、TaskExecutor和自动启动
接下来看容器配置,这是整个接入过程中容易写错的地方。一个标准的RedisMessageListenerContainer配置如下:
@Configuration public class RedisPubSubConfig { @Bean public RedisMessageListenerContainer redisMessageListenerContainer( RedisConnectionFactory connectionFactory, CacheInvalidateListener cacheInvalidateListener) { RedisMessageListenerContainer container = new RedisMessageListenerContainer(); container.setConnectionFactory(connectionFactory); container.setTopicSerializer(new StringRedisSerializer()); // 使用有界线程池来执行订阅任务 ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); executor.setCorePoolSize(4); executor.setMaxPoolSize(16); executor.setQueueCapacity(20); executor.setThreadNamePrefix("redis-sub-task-"); executor.initialize(); container.setTaskExecutor(executor); container.addMessageListener(cacheInvalidateListener, new ChannelTopic("cache:invalidate")); return container; } }这段配置里有三个关键点。第一,setTopicSerializer一定要设成StringRedisSerializer,原因就是要保证订阅的频道字节和发布端的UTF-8字节完全一致,这个细节省掉能避免一个极其隐蔽的"收不到消息"问题,后面坑里细讲。第二,手动指定了TaskExecutor,这是因为容器底层每个订阅任务都是阻塞线程,默认的SimpleAsyncTaskExecutor会在每次任务到来时直接新建线程,订阅多了线程数不可控。第三,容器作为Bean返回后,Spring的Lifecycle处理器会自动调用它的start()方法,不需要手动启动。
4.4 发送端:convertAndSend与序列化策略
发送端更简单,把StringRedisTemplate注入到你的业务类,一行调用即可:
@Service public class ProductService { private final StringRedisTemplate stringRedisTemplate; public void updateProduct(Long productId) { // 1. 更新数据库 productDao.updatePrice(productId); // 2. 发出缓存失效广播 stringRedisTemplate.convertAndSend("cache:invalidate", "product:" + productId); } }convertAndSend方法做的事情,本质就是把channel和message都用StringRedisSerializer编码成UTF-8字节,然后通过PUBLISH命令发给Redis。这里有一个我一直强调的原则:发送端和接收端的序列化策略必须完全对称。如果你的发送端用的是默认RedisTemplate,默认序列化器是JdkSerializationRedisSerializer,那么消息体是一堆以\xAC\xED开头的JDK二进制垃圾,接收端的UTF-8转String就是乱码。这就是为什么我整个Demo都围绕StringRedisTemplate展开。
如果团队项目里已经约定使用JSON格式来传递消息体,建议发送的时候自己先把对象转成JSON字符串再调convertAndSend,接收端用Jackson从字符串反序列化。这样既保证跨语言可读性,又把序列化边界控制在业务代码内,是最透明的方案。
4.5 验证步骤与预期日志
代码写完后,验证流程大概是这样的:先保证Redis已经启动,然后正常启动Spring Boot应用。你会看到控制台日志里出现类似Initializing RedisMessageListenerContainer和Started redis-sub-task-1之类的信息,这说明容器已经启动并进入订阅状态。接着在redis-cli里执行:
redis-cli publish cache:invalidate "product:1001"如果一切正常,应用控制台会打印出:
收到缓存失效广播,channel=cache:invalidate, key=product:1001注意这条日志出现的同时,没有报任何序列化异常,说明整个链路通了。想验证多实例效果,就同时启动两个应用实例,分别监听同一个频道,再用redis-cli发一条消息,你会看到两个实例都收到了同一条消息。这正是广播的核心语义:一条消息,全部订阅者各收一份。
5. 进阶用法:模式订阅、动态订阅、多监听器的组合编排
5.1 PatternTopic:一个监听器接收一类频道
很多场景下你不会只关心固定的一个频道,而是一组频道。比如订单域的消息,有order:created、order:paid、order:cancelled,你希望一个监听器统一处理,避免为每个频道重复注册一段几乎一样的代码。Redis的模式订阅就是干这个的,Spring Data Redis对应的实现是PatternTopic。
container.addMessageListener( orderMessageListener, new PatternTopic("order:*") );使用模式订阅后,onMessage回调里第二个参数pattern就不再是空值,而是这条消息实际匹配到的模式。比如一个事件发到order:paid频道,监听器收到的pattern就是order:*的字节。你可以利用这个参数做二次路由,也可以完全忽略它。
这里要提醒一个坑:Redis的模式匹配通配符支持*、?、[]字符集,但模式匹配发生在Redis服务端,不是在客户端。也就是说"pattern匹配到哪些频道"这件事由Redis决定,你无法在Spring容器层面干预。如果你发出了order:paid:2024这个消息,而你的pattern是order:*,消息照样会被收下来,频道粒度全由你的事件命名规范来控制。所以事件命名的时候一定要定义好层级分隔符,别把领域信息混在格式里,否则后面做模式匹配会很痛苦。
5.2 动态注册与移除:容器运行时的API
容器的另一个实用价值,是允许在运行期动态调整订阅关系。注册监听器不止可以在@Bean配置里做,你也可以把它封装为一个可管理组件,把"是否监听某个频道"变成一种可配置的能力:
@Service public class DynamicSubscriptionService { private final RedisMessageListenerContainer container; private final Map<String, MessageListener> listeners = new ConcurrentHashMap<>(); public void subscribe(String channel, MessageListener listener) { listeners.put(channel, listener); container.addMessageListener(listener, new ChannelTopic(channel)); } public void unsubscribe(String channel) { MessageListener listener = listeners.remove(channel); if (listener != null) { container.removeMessageListener(listener); } } }这个能力在什么场景下有用?我实际遇到的一个需求是多租户系统。每个租户有自己独立的一批事件频道,租户开通时动态订阅,租户停用时动态移除,不需要重启应用。另外一个场景是运行时配置中心推送了新的"监听开关",你可以在不停机的情况下热切换。
但要注意:动态添加监听器时,容器会在内部为这个新的订阅任务分配一个新的阻塞线程和一个新的Redis连接。如果频繁增删,又没有复用合理的线程池,连接和线程会产生较大波动。所以动态订阅虽然灵活,也应该维持在一个可控的数量级,不要像写JVM内存缓存那样随手往里塞。
5.3 多监听器的职责拆分与线程隔离
一个RedisMessageListenerContainer可以注册任意多个监听器,对应任意多个Topic。当监听器多了,这些监听器之间共享的是同一个TaskExecutor线程池和同一个Dispatcher调度逻辑。如果你的系统里有两类消息,一类是业务核心的订单事件,另一类是低优级的日志采集,把它们塞进同一个容器不是不可以,但要接受一个问题:订阅任务之间的线程分配是相对独立的(每个topic都有一个阻塞订阅线程),但消息派发到监听器这一段的执行是共享线程的。高优级的监听器代码如果执行很慢,会把任务线程池占满,拖累其他监听器的消息处理。
所以我的经验是:按消息重要程度拆分容器。订单事件一个container,给它一个核心线程多一点的线程池;日志采集另一个container,用一个保守的线程池。两个容器各自启动、各自管理连接,故障隔离做得更好。代价只是多注册几个Bean,完全值得。
6. 实踩过的坑:序列化、线程模型、连接恢复与消息丢失
6.1 SerializationException:最典型的序列化问题排查链路
我最早接入Redis Pub/Sub时,犯过一个教科书式的错误。为了图省事,用默认RedisTemplate做发送端,监听器用了Spring的MessageListenerAdapter,并声明了一个接收String参数的方法。结果一跑,控制台直接抛异常,核心内容是:
java.lang.IllegalStateException: Cannot deserialize; nested exception is org.springframework.core.serializer.support.SerializationFailedException: ...当时第一反应是怀疑接收端配置有问题,查了半天容器配置也没看出毛病。后来用redis-cli手动publish了一条带中文的消息,再用redis-cli monitor观察,发现Redis里存的消息内容根本不是可读文本,而是以\xAC\xED\x00\x05开头的一串二进制。这才意识到问题出在发送端:默认RedisTemplate的valueSerializer是JdkSerializationRedisSerializer,它把字符串当Java对象做JDK序列化了,接收端却天真地按UTF-8处理。
这个坑的排查思路值得总结一下。遇到序列化异常,不要先改接收端代码,而是先确认发到Redis的消息原始字节是什么样。用redis-cli monitor观察,或者干脆用可视化管理工具连上去,找到那个频道,看清消息是纯文本还是二进制乱码。发送端用什么序列化器,接收端就必须用什么序列化器,这是铁律。我的最终方案就是把发送端换成StringRedisTemplate,收发两端全部走UTF-8字节,之后再没出现过这个问题。
另外提一句,如果你非要用MessageListenerAdapter省掉手动转String这一步,那么你要同时保证容器使用的消息转换器和发送端序列化器一致。这需要设置adapter的Serializer,代码会多出好几行,收益却只是省了一个new String,我个人认为不划算。
6.2 频道字节不一致:消息发布成功但监听器无反应的根源
这个坑更隐蔽,因为它不报任何异常。现象就是:redis-cli publish消息,提示返回了1,说明确实有一个订阅者在线;但应用日志里死活没有输出。我排查了很久,最后发现是频道字节不一致的问题。
前面说过,RedisMessageListenerContainer在调用SUBSCRIBE命令时,需要把频道名序列化成字节。如果容器没有设置TopicSerializer,它默认使用的是JDK序列化器。也就是说,我订阅的频道名"cache:invalidate"被编码成的字节,和StringRedisTemplate发送端用UTF-8编码的频道名完全不是同一个东西。Redis服务器以为你订阅的是"A频道",发布的却是"B频道",两边频道不一样,消息当然投递不到。
排查过程也很有意思。我先在Redis服务器上执行PUBSUB CHANNELS命令,列出当前所有被订阅的活跃频道,发现列表里的频道名后面跟着一堆不可见字符,再比对监控端的订阅信息,才定位到问题。解决方案就是配置里那行container.setTopicSerializer(new StringRedisSerializer())。这个细节太容易漏了,网上很多教程都不提,导致不少人照着抄也会莫名踩坑。
6.3 线程数暴涨:阻塞订阅模型下的资源规划
再讲一个生产环境遇到的性能问题。我们的监听器数量不多,只有五六个ChannelTopic,但应用启动后观察监控面板,发现Redis连接数比预期的多出不少,线程数也在慢慢爬升。后来用jstack看线程快照,看到大量redis-sub-task-前缀的线程,才意识到订阅模型和普通请求处理模型是不同的。
RedisMessageListenerContainer的内部分工是:每个Topic对应一个Subscription任务,这个任务是阻塞式的,需要一条独占的Redis连接和一个独占的执行线程。所以一个容器里监听器的注册数量直接影响的就是线程数量和连接数量。默认情况下,传给容器的TaskExecutor如果不指定,Spring Data Redis使用的是SimpleAsyncTaskExecutor,它的特点是来一个任务就new一个线程,完全没有复用。
解决思路有两个层面。第一,显式配置线程池,就像我在Demo里做的那样,用ThreadPoolTaskExecutor把线程数量约束住。注意这个线程池的线程可能长期处于阻塞状态,所以corePoolSize要设定为和你预期订阅总量相当的数量。第二,控制监听器的粒度,看看能否用PatternTopic把一类频道合并成一个监听器订阅,而不是每个频道都单独注册一次。这两个措施配合下来,线程和连接的问题基本可以缓解。
6.4 连接断开与自动重连:主从切换时的订阅恢复
Redis连接不会永远稳定,网络抖动、Redis主从切换、应用与Redis之间的防火墙空闲超时,都会导致订阅连接被中断。好消息是RedisMessageListenerContainer内部有重连机制,它会在连接异常后自动重建连接、重新执行订阅。这不是一个需要你额外处理的故障,但有两个细节值得了解。
第一个细节:连接中断后,业务消息的丢失窗口是必然存在的。因为Pub/Sub没有持久化,断开期间发布者发出来的消息就是丢了,重连成功后也只从重连时刻开始收新消息。对于缓存失效这类场景可以容忍,但如果你在用它做在线状态变更通知,客户端这边就要有兜底逻辑。第二个细节:主从切换时,即使你使用的是主节点的连接地址,切换过程中连接也会被重置。应用层通常不会有感知,因为容器会自动重连到新主节点。但如果你的Redis架构是读写分离,且不小心把订阅连接配置到了从节点上,遇到主从切换或版本语义变化时,订阅的恢复行为就可能不符合预期。所以我的建议是订阅专用的连接工厂尽量直连主节点,别在订阅通道上玩读写分离。
6.5 消息丢失的补偿哲学:Pub/Sub不是可靠的投递通道
最后这个"坑"其实是机制特性,不算Bug。我要特别强调:Redis Pub/Sub没有消息确认,没有重试机制,发布者发出消息后根本不知道有多少订阅者收到了,甚至不知道有没有订阅者在听。这是它最轻量、最实时的代价。如果你用它做重要业务事件的传递,一旦某个实例正在重启、GC长暂停、或者网络闪断,这条消息就永久丢失了。
正确的防御姿势是设计兜底。比如做缓存失效广播时,缓存本身要保留合理的TTL,让最坏情况下数据也能在TTL后自动恢复一致。做用户状态更新时,除了广播通知,还要保留一个定时任务做状态对账。这些都是Pub/Sub使用哲学的一部分:广播负责快,兜底负责稳。千万别把"发出去"当成"送达并成功处理",否则早晚会被线上事故教育。
7. 边界判断:什么时候该换Redis Streams或消息队列
7.1 Redis Streams补齐了Pub/Sub哪三块短板
如果你发现业务需求逐渐超出了Pub/Sub的能力边界,可以考虑先用Redis Streams过渡,而不是一步跳到导入RabbitMQ或Kafka。Redis Streams是Redis 5.0引入的数据类型,它和Pub/Sub的核心区别是:消息会被持久化保存在Redis里,消费者可以按需读取,未被消费的历史消息保留在流中,并且支持消费者组和ACK确认机制。
它对Pub/Sub的补强体现在三块:一是消息不丢,消费者离线期间的消息保留在Stream里,回来后可以继续读取;二是消费进度可控,消费者处理成功后XACK确认,失败的消息可以重新进入处理流程;三是支持消费组而不是纯广播,一条消息只被组内一个成员消费,这就能做任务分发。如果你的需求是要"消息不丢、能重试、能分组",但暂时不想引入重量级MQ,Streams是合理中间站。
Spring Data Redis 3.x对Streams的支持也已经很完善,有StreamMessageListenerContainer,有StreamOperations供RedisTemplate操作,同样可以像Pub/Sub那样注册监听器。不过Streams的模型比Pub/Sub复杂不少,涉及Stream、Group、Consumer、Pending Entries List等概念,学习曲线陡一些。
7.2 什么时候必须换MQ:三点硬性信号
我的经验是,当项目出现下面三个信号之一时,就不要再犹豫,直接上MQ:
第一个信号,消息处理要求可靠投递,丢一条消息会造成资金差错或数据不一致。这不是靠兜底能掩盖的,需要MQ的持久化和回执机制。第二个信号,存在生产速率和消费速率不匹配的情况,需要消息堆积能力。Pub/Sub根本堆积不了,Redis Streams堆积过多也会撑爆内存,但MQ天生就是为削峰填谷设计的。第三个信号,你需要复杂的路由规则、延迟消息、死信队列这些高级特性。这些在Redis世界里要靠自己造轮子,而主流MQ开箱即用。
这三个信号出现任何一个,Redis Pub/Sub就得退场。我曾经在一个项目中硬是把Pub/Sub用于订单事件通知,结果高峰期消息洪水把下游打了个措手不及,后面沉寂下来老老实实改造,整个过程非常痛苦。
7.3 我的选型经验:先广播后加固的演进路线
这几年的实践让我形成了一个比较稳定的选型路线:新项目第一版,凡是广播通知类的需求,先考虑Redis Pub/Sub,因为现有基础设施不用额外投入,代码简单,业务方也能快速理解。当消息量级增长、可靠性要求提高之后,再针对确有需求的场景逐步升级到Redis Streams,或者干脆引入MQ。我做缓存失效广播的项目跑了两年,直到现在还是Redis Pub/Sub,因为它一直处在适用边界内。
这里想再分享一个我个人的实践惯,就是所有频道的命名一定要带上环境前缀和环境标识,比如prod:cache:invalidate、dev:cache:invalidate。因为Redis Pub/Sub是全库广播,两个环境共用一套Redis时,频道名如果不区分,开发环境的缓存失效消息就会跑到生产实例里被处理。这种问题定位起来极其费劲,但一条命名规范就能规避。另外一个实用小技巧是多实例部署时,可以在消息体里带上实例ID,接收端判断是不是自己发出去的消息,需要时直接忽略,避免"自己通知自己"形成无意义的回环。这些都是文档里不会写、但真实项目里每天都会面对的小事。