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

Spring订阅Redis集群Key过期事件失败求助

环境信息

  • Mac Ventura 13.6.3
  • Temurin 17
  • SpringBoot 3.2.1
  • 依赖:org.springframework.boot:spring-boot-starter-data-redis
  • 本地Redis集群(localhost:7001、localhost:7002、localhost:7003)

需求

实现通过Spring接收Redis集群的Key过期事件。

已执行操作

连接集群后执行以下Redis命令:

CONFIG SET notify-keyspace-events Ex

SET hi 123
EXPIRE hi 3

Java代码

配置类

package kr.co.yogiyo.payo.infrastructure.temporal;

import io.lettuce.core.ReadFrom;
import kr.co.yogiyo.payo.infrastructure.temporal.service.RedisExpireEventService;
import lombok.RequiredArgsConstructor;
import org.springframework.boot.autoconfigure.data.redis.RedisConnectionDetails;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.data.redis.connection.RedisClusterConfiguration;
import org.springframework.data.redis.connection.RedisConnectionFactory;
import org.springframework.data.redis.connection.RedisNode;
import org.springframework.data.redis.connection.lettuce.LettuceClientConfiguration;
import org.springframework.data.redis.connection.lettuce.LettuceConnectionFactory;
import org.springframework.data.redis.listener.PatternTopic;
import org.springframework.data.redis.listener.RedisMessageListenerContainer;

@RequiredArgsConstructor
@Configuration
public class RedisConfig {

  private final RedisConnectionDetails redisConnectionDetails;

  @Bean
  public RedisConnectionFactory redisConnectionFactory() {
    RedisConnectionDetails.Cluster cluster = redisConnectionDetails.getCluster();
    RedisClusterConfiguration clusterConfiguration = new RedisClusterConfiguration();
    for (RedisConnectionDetails.Node node : cluster.getNodes()) {
      clusterConfiguration.addClusterNode(new RedisNode(node.host(), node.port()));
    }

    LettuceClientConfiguration clientConfig =
        LettuceClientConfiguration.builder().readFrom(ReadFrom.REPLICA_PREFERRED).build();

    return new LettuceConnectionFactory(clusterConfiguration, clientConfig);
  }

  @Bean
  RedisMessageListenerContainer keyExpirationListenerContainer(
      RedisConnectionFactory connectionFactory, RedisExpireEventService redisExpireEventService) {
    RedisMessageListenerContainer listenerContainer = new RedisMessageListenerContainer();
    listenerContainer.setConnectionFactory(connectionFactory);
    listenerContainer.addMessageListener(
        redisExpireEventService, new PatternTopic("__keyevent@*__:expired"));
    return listenerContainer;
  }
}

服务类

package kr.co.yogiyo.payo.infrastructure.temporal.service;

import lombok.RequiredArgsConstructor;
import org.springframework.data.redis.connection.Message;
import org.springframework.data.redis.connection.MessageListener;
import org.springframework.stereotype.Service;

@Service
@RequiredArgsConstructor
public class RedisExpireEventService implements MessageListener {
  @Override
  public void onMessage(Message message, byte[] pattern) {
    System.out.println("Message received: " + message.toString());
  }
}

application.yml

spring:
  main:
    banner-mode: off
    allow-bean-definition-overriding: true
  jackson:
    property-naming-strategy: SNAKE_CASE
  data:
    redis:
      host: ${PAYO_REDIS_MASTER_HOST:localhost}
      port: ${PAYO_REDIS_MASTER_PORT:7002}
      cluster:
        nodes: ${PAYO_REDIS_REPLICATION_NODES:localhost:7001,localhost:7002,localhost:7003}

问题

Spring项目可正常启动,但Key过期时无任何响应(调试模式下断点也未触发)。

可正常运行的Python代码

import redis


def main():
    # 连接Redis
    r = redis.Redis(host='localhost', port=7002, db=0)

    # 订阅__keyevent@0__:expired频道
    pubsub = r.pubsub()
    pubsub.psubscribe('__keyevent@*__:expired')

    # 开始监听消息
    for message in pubsub.listen():
        print("有键过期! ", message)

if __name__ == "__main__":
    main()

运行输出

有键过期!  {'type': 'psubscribe', 'pattern': None, 'channel': b'__keyevent@*__:expired', 'data': 1}
有键过期!  {'type': 'pmessage', 'pattern': b'__keyevent@*__:expired', 'channel': b'__keyevent@0__:expired', 'data': b'hi'}

提问

我的Java代码中缺少了什么配置或逻辑?


问题排查与解决方案

  1. 全节点配置notify-keyspace-events
    Redis集群中每个节点独立运行,你只在单个节点执行了CONFIG SET notify-keyspace-events Ex,但Key可能存储在任意节点,过期事件只会在Key所在节点触发。需要给集群所有节点执行该命令;若要永久生效,需修改每个节点的redis.conf文件,添加notify-keyspace-events Ex后重启节点。

  2. 配置Lettuce集群拓扑刷新
    集群模式下,Lettuce默认仅在初始连接节点订阅频道,无法接收其他节点的过期事件。需添加拓扑刷新配置,让客户端自动发现集群节点并全节点订阅:
    修改redisConnectionFactory方法中的Lettuce客户端配置:

    import io.lettuce.core.ClientOptions;
    import io.lettuce.core.cluster.ClusterTopologyRefreshOptions;
    
    // ...
    LettuceClientConfiguration clientConfig = LettuceClientConfiguration.builder()
        .readFrom(ReadFrom.REPLICA_PREFERRED)
        .clientOptions(ClientOptions.builder()
            .clusterTopologyRefreshOptions(ClusterTopologyRefreshOptions.builder()
                .enableAllAdaptiveRefreshTriggers()
                .build())
            .build())
        .build();
    
  3. 指定Redis数据库编号
    Python代码连接的是db=0,检查Spring是否使用相同数据库,若默认不是db=0,需在application.yml中添加:

    spring:
      data:
        redis:
          database: 0
    
  4. 订阅专用连接工厂(可选)
    过期事件仅在主节点触发,REPLICA_PREFERRED策略可能导致订阅连接到副本节点无法接收事件。可给订阅容器单独配置主节点优先的连接工厂:

    @Bean
    RedisConnectionFactory redisSubscriptionConnectionFactory() {
      RedisConnectionDetails.Cluster cluster = redisConnectionDetails.getCluster();
      RedisClusterConfiguration clusterConfiguration = new RedisClusterConfiguration();
      for (RedisConnectionDetails.Node node : cluster.getNodes()) {
        clusterConfiguration.addClusterNode(new RedisNode(node.host(), node.port()));
      }
    
      LettuceClientConfiguration clientConfig = LettuceClientConfiguration.builder()
          .readFrom(ReadFrom.MASTER)
          .clientOptions(ClientOptions.builder()
              .clusterTopologyRefreshOptions(ClusterTopologyRefreshOptions.builder()
                  .enableAllAdaptiveRefreshTriggers()
                  .build())
              .build())
          .build();
    
      return new LettuceConnectionFactory(clusterConfiguration, clientConfig);
    }
    
    @Bean
    RedisMessageListenerContainer keyExpirationListenerContainer(
        RedisExpireEventService redisExpireEventService) {
      RedisMessageListenerContainer listenerContainer = new RedisMessageListenerContainer();
      listenerContainer.setConnectionFactory(redisSubscriptionConnectionFactory());
      listenerContainer.addMessageListener(
          redisExpireEventService, new PatternTopic("__keyevent@*__:expired"));
      return listenerContainer;
    }
    

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 09:07:05