因为pubsub.listen()是阻塞式迭代器,会卡住tornado单线程ioloop,导致无法处理http请求等异步任务;正确做法是用aioredis v2.0.1配合ioloop.add_reader手动监听socket并parse_response。

为什么不能直接在 Tornado 的 IOLoop 中用 redis-py 的 pubsub.listen()
因为 pubsub.listen() 是一个阻塞式迭代器,它内部调用 connection.read_response() 并持续轮询 socket,会彻底卡住 Tornado 的单线程事件循环。一旦启动,IOLoop 就无法处理 HTTP 请求、定时任务或其他回调。
常见错误现象:curl http://localhost:8000/health 无响应,日志停在 Starting pubsub listener... 后再无输出;Tornado 进程 CPU 占用低但完全不响应新请求。
- redis-py 的
PubSub默认基于同步 socket,和 Tornado 的异步模型不兼容 - 即使把
listen()放进run_in_executor,消息到达时也无法安全触发RequestHandler或修改Application状态(涉及线程安全与 IOLoop 绑定) - 不要尝试用
setblocking(False)+select手动轮询 —— redis-py 内部缓冲和协议解析逻辑不暴露给用户,不可靠
正确做法:用 aioredis + IOLoop.add_reader 手动接管 socket
aioredis(v2.x)提供原生 asyncio 支持,其底层连接对象暴露了可等待的 socket 文件描述符。我们可以把它“嫁接”进 Tornado 的 IOLoop,让事件循环主动通知有数据可读,而非让 Redis 客户端自己阻塞等待。
关键点在于:不调用 pubsub.listen(),而是用 pubsub.execute_command("SUBSCRIBE", ...) 发起订阅,然后监听 socket 可读事件,手动调用 pubsub.parse_response() 处理响应。
- 必须使用
aioredis==2.0.1(非 1.x 或 3.x+),因 v2 是唯一同时支持 asyncio 和暴露connection._sock的稳定版本 - 订阅后需禁用
pubsub.auto_reconnect = False,否则底层重连逻辑会破坏手动 socket 管理 - 每次
parse_response()成功后,要立即再次调用IOLoop.current().add_reader(...),因为 Tornado 的 reader 只触发一次 - 示例片段:
import aioredis from tornado.ioloop import IOLoop <p>async def init_pubsub(): redis = await aioredis.from_url("redis://localhost") pubsub = redis.pubsub() await pubsub.subscribe("channel:a") conn = pubsub.connection</p><h1>确保连接已建立且 socket 可读</h1><pre class="brush:php;toolbar:false;">loop = IOLoop.current() def on_message(): try: msg = pubsub.parse_response(block=False) # 非阻塞 if msg and msg[0] == b'message': print("Got:", msg[2]) except (ConnectionError, OSError): pass # 断开时由上层重连逻辑处理 finally: # 重新注册,保持监听 loop.add_reader(conn._sock, on_message) loop.add_reader(conn._sock, on_message)如何安全地把 Redis 消息转发给 Web 请求上下文
Tornado 中没有全局“会话池”或“客户端广播总线”,消息到达后若想推送给特定 HTTP 连接(比如 WebSocket),必须自行维护活跃连接引用,并确保操作发生在 IOLoop 线程内。
- 避免在
on_message回调里直接调用websocket.write_message()—— 如果该 WebSocket 已关闭,会抛StreamClosedError,且未被 try/catch 包裹时将终止整个 reader 回调链 - 推荐模式:把消息丢进一个
queue.Queue(注意是线程安全的queue.Queue,不是 asyncio.Queue),再用IOLoop.current().spawn_callback()异步消费,这样异常不会影响 socket 监听 - WebSocket 连接需在
open()时加入全局集合(如set),在on_close()或check_origin=False失败时及时移除,否则内存泄漏且消息误投 - 不要依赖
self.application存全局 pubsub 实例 —— Tornado 应用实例是线程局部的,而 pubsub socket 监听在主线程,没问题;但若未来改用多进程部署,需改用 Redis Stream 或外部消息队列解耦
生产环境必须加的兜底机制
真实场景中,Redis 连接闪断、订阅丢失、IOLoop 调度延迟都可能发生。纯靠
add_reader不足以维持稳定监听。- 必须实现心跳检测:每 30 秒发一次
PING,收不到PONG则主动断开并重建 pubsub 连接 - 订阅命令失败(如返回
None或抛ReplyError)时,不能静默忽略 —— 要记录日志并退避重试(例如指数退避 1s → 2s → 4s) - 在
on_message中对parse_response()做超时控制:传入timeout=0.1,防止因 Redis 协议解析 bug 卡死 - 如果项目已用
tornado.web.Application.settings["redis_pool"]管理连接,pubsub 必须独占一个连接(minsize=maxsize=1),否则其他协程可能意外关闭该 socket
最易被忽略的是连接复用边界:HTTP handler 用的 redis client 和 pubsub 用的 client 必须物理隔离,哪怕它们连的是同一个 Redis 实例 —— 否则
client.close()可能提前终结 pubsub socket。 - 避免在
Python免费学习笔记(深入):立即使用
在学习笔记中,你将探索 Python 的核心概念和高级技巧!











