
Redis发布订阅在高并发下出现延迟,根本原因不是Redis发不出消息,而是消费端卡在了pubsub.listen()的阻塞循环里——它单线程、不批处理、不做超时控制,一卡全卡。
Python中用多线程启动listen()必须隔离连接
每个订阅线程必须持有独立的redis.Redis实例和pubsub对象,不能共用一个连接或连接池。共用会导致:
-
ConnectionError: Connection closed by server—— 因为其他线程调用unsubscribe()或reset()会重置整个连接状态 - 消息错收:A线程订阅
order:paid,B线程误发PUBLISH order:canceled ...,服务端可能直接断连 -
client_longest_output_list持续上涨,超过1000就说明消费者已严重积压
正确写法是在线程内新建连接:
Redis 缓存和数据结构管理技能。通过自然语言操作 Redis,支持 String、Hash、List、Set、ZSet、Stream 等数据结构操作。当用户提到 Redis、缓存、消息队列、会话存储时使用此技能。
def subscribe_to_channel(channel_name):
r = redis.Redis(host='127.0.0.1', port=6379, decode_responses=True, socket_timeout=5)
pubsub = r.pubsub()
pubsub.subscribe(channel_name)
for msg in pubsub.listen():
if msg['type'] == 'message':
process_message(msg['data']) # 真实处理逻辑
阻塞式listen()里不能做耗时操作
pubsub.listen()是同步阻塞迭代器,一旦你在循环里调用数据库写入、HTTP请求或复杂JSON解析,下一条消息就得等几秒甚至几十秒,缓冲区立刻堆积。
- 把反序列化、路由判断、日志记录等逻辑移到异步任务队列(如
concurrent.futures.ThreadPoolExecutor)中执行 - 设置
socket_timeout=3,避免网络抖动导致无限挂起 - 加
try/except RedisError捕获断连,自动重连并重新subscribe
Java用Lettuce时别共享StatefulRedisPubSubConnection
Lettuce的StatefulRedisPubSubConnection是线程不安全的,多个业务模块共用一个连接对象,极易触发READONLY You can't write against a read only replica错误——这其实不是权限问题,而是连接被其他模块误发了PUBLISH,导致内部状态混乱,服务端强制断开。
- 每个订阅逻辑(如库存变更监听、支付回调监听)应使用独立的
ClientResources和EventLoopGroup - 配置
spring.redis.lettuce.pool.max-active=200只对命令连接有效,Pub/Sub连接不走这个池,必须显式创建新连接 - 用
connection.addListener()注册回调,而非长期阻塞sync().listen()
真正卡住的从来不是Redis,而是你没意识到listen()是个单点瓶颈,且默认行为完全不适应高并发消费场景——它不批、不异步、不熔断,只安静地等你把它拖垮。










