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

高并发事件下Redis键空间通知遇Unexpected end of stream错误

Redis键空间通知在高并发过期场景下的连接异常问题

我正在探索Redis键空间通知及其在SpringBoot中的实现,用于TTL过期场景。测试承载规模时,1万-10万键同时过期的情况下运行正常,但当同一时间过期的键增至50万时,出现如下连接异常:

org.springframework.data.redis.RedisConnectionFailureException: Unexpected end of stream.; nested exception is redis.clients.jedis.exceptions.JedisConnectionException: Unexpected end of stream.
at org.springframework.data.redis.connection.jedis.JedisExceptionConverter.convert(JedisExceptionConverter.java:65) ~[spring-data-redis-2.7.5.jar:2.7.5]
at org.springframework.data.redis.connection.jedis.JedisExceptionConverter.convert(JedisExceptionConverter.java:42) ~[spring-data-redis-2.7.5.jar:2.7.5]
at org.springframework.data.redis.PassThroughExceptionTranslationStrategy.translate(PassThroughExceptionTranslationStrategy.java:44) ~[spring-data-redis-2.7.5.jar:2.7.5]
at org.springframework.data.redis.FallbackExceptionTranslationStrategy.translate(FallbackExceptionTranslationStrategy.java:42) ~[spring-data-redis-2.7.5.jar:2.7.5]
at org.springframework.data.redis.connection.jedis.JedisConnection.convertJedisAccessException(JedisConnection.java:192) ~[spring-data-redis-2.7.5.jar:2.7.5]
at org.springframework.data.redis.connection.jedis.JedisConnection.doWithJedis(JedisConnection.java:836) ~[spring-data-redis-2.7.5.jar:2.7.5]
at org.springframework.data.redis.connection.jedis.JedisConnection.pSubscribe(JedisConnection.java:741) ~[spring-data-redis-2.7.5.jar:2.7.5]
at org.springframework.data.redis.listener.RedisMessageListenerContainer$Subscriber.doSubscribe(RedisMessageListenerContainer.java:1231) ~[spring-data-redis-2.7.5.jar:2.7.5]
at org.springframework.data.redis.listener.RedisMessageListenerContainer$BlockingSubscriber.lambda$eventuallyPerformSubscription$2(RedisMessageListenerContainer.java:1433) ~[spring-data-redis-2.7.5.jar:2.7.5]
at java.base/java.lang.Thread.run(Thread.java:833) ~[na:na]
Caused by: redis.clients.jedis.exceptions.JedisConnectionException: Unexpected end of stream.
at redis.clients.jedis.util.RedisInputStream.ensureFill(RedisInputStream.java:202) ~[jedis-3.8.0.jar:na]
at redis.clients.jedis.util.RedisInputStream.readLongCrLf(RedisInputStream.java:154) ~[jedis-3.8.0.jar:na]
at redis.clients.jedis.util.RedisInputStream.readIntCrLf(RedisInputStream.java:148) ~[jedis-3.8.0.jar:na]
at redis.clients.jedis.Protocol.processBulkReply(Protocol.java:188) ~[jedis-3.8.0.jar:na]
at redis.clients.jedis.Protocol.process(Protocol.java:170) ~[jedis-3.8.0.jar:na]
at redis.clients.jedis.Protocol.processMultiBulkReply(Protocol.java:221) ~[jedis-3.8.0.jar:na]
at redis.clients.jedis.Protocol.process(Protocol.java:172) ~[jedis-3.8.0.jar:na]
at redis.clients.jedis.Protocol.read(Protocol.java:230) ~[jedis-3.8.0.jar:na]
at redis.clients.jedis.Connection.readProtocolWithCheckingBroken(Connection.java:352) ~[jedis-3.8.0.jar:na]
at redis.clients.jedis.Connection.getUnflushedObjectMultiBulkReply(Connection.java:314) ~[jedis-3.8.0.jar:na]
at redis.clients.jedis.BinaryJedisPubSub.process(BinaryJedisPubSub.java:101) ~[jedis-3.8.0.jar:na]
at redis.clients.jedis.BinaryJedisPubSub.proceedWithPatterns(BinaryJedisPubSub.java:89) ~[jedis-3.8.0.jar:na]
at redis.clients.jedis.BinaryJedis.psubscribe(BinaryJedis.java:3900) ~[jedis-3.8.0.jar:na]
at org.springframework.data.redis.connection.jedis.JedisConnection.lambda$pSubscribe$7(JedisConnection.java:746) ~[spring-data-redis-2.7.5.jar:2.7.5]
at org.springframework.data.redis.connection.jedis.JedisConnection.doWithJedis(JedisConnection.java:834) ~[spring-data-redis-2.7.5.jar:2.7.5]
... 4 common frames omitted

实现代码

过期键监听器

package com.example.rediseventlistnertest.listner;

import lombok.extern.slf4j.Slf4j;
import org.springframework.data.redis.connection.Message;
import org.springframework.data.redis.connection.MessageListener;
import org.springframework.stereotype.Component;

import java.util.ArrayList;
import java.util.Date;
import java.util.List;

@Slf4j
@Component
public class EventTimeOutListner implements MessageListener {

    public List<String> keys = new ArrayList<>();

    @Override
    public void onMessage(Message message, byte[] pattern) {
        String key = new String(message.getBody());
        Date date = new Date();
        Long sec = date.getTime();
        String d = date.toString();
        keys.add(key);
        log.info("expired key: {}, expired at {} {} {}", key, sec, date, keys.size());
    }
}

Redis监听配置类

@Slf4j
@Configuration
public class RedisConfiguration {

    @Bean
    public RedisMessageListenerContainer redisMessageListenerContainer(RedisConnectionFactory redisConnectionFactory, EventTimeOutListner taskTimeoutListener, TimeOutListner tol) {

        RedisMessageListenerContainer listenerContainer = new RedisMessageListenerContainer();
        listenerContainer.setMaxSubscriptionRegistrationWaitingTime(10000);
        try {
            listenerContainer.setConnectionFactory(redisConnectionFactory);
            listenerContainer.addMessageListener(taskTimeoutListener,
                    new PatternTopic("__keyevent@*__:expired"));
            log.info("size of  keys : {}",  taskTimeoutListener.keys.size());
            //listenerContainer.addMessageListener(tol,
            //        new PatternTopic("__keyevent@*__:expired"));
            listenerContainer.setErrorHandler(
                    e -> log.error("Error in redisMessageListenerContainer", e));
        } catch (Exception e) {
            log.info("Something went worng exception is : {}", e.getMessage());
            e.printStackTrace();
        } finally {
            return listenerContainer;
        }
    }
}

问题与求助

我尝试过添加连接池、增加连接数等方法,但问题仍未解决。无法确定确切原因,想请教:

  • 有没有人遇到过类似问题?
  • 可能的原因是什么?
  • 如何避免该故障并扩展Redis以处理大规模键空间通知?

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.14 05:40:32