不能直接用 aioredis.from_url() 创建 pubsub 实例,因为 redis 实例设计为短时命令执行,连接会被复用或关闭;而 pubsub 需长期挂起等待消息,必须单独调用 redis.pubsub() 获取专用长连接对象。

为什么不能直接用 aioredis.from_url() 创建 pubsub 实例
因为 aioredis.Redis 实例本身不支持长期保持 pubsub 连接——它被设计为短时命令执行,底层连接会在空闲后被复用或关闭。而 pubsub 要求连接持续挂起等待消息,否则会断连丢消息。
正确做法是显式创建 aioredis.Redis 实例用于发消息(publish),再单独用 redis.pubsub() 获取一个专用的 PubSub 对象来监听(subscribe):
-
redis = aioredis.from_url("redis://localhost")→ 仅用于publish -
pubsub = redis.pubsub()→ 必须调用此方法获取独立的长连接对象 - 这个
pubsub实例必须自己管理生命周期,不能和普通 redis 客户端混用连接池
如何避免 pubsub.get_message() 阻塞主线程
pubsub.get_message(block=True) 是同步阻塞调用,在协程里直接用会导致整个 event loop 卡住;而 block=False 又容易轮询浪费 CPU。
推荐方案是用 pubsub.listen() 配合 async for ——它内部已封装为异步迭代器,能自然挂起等待新消息:
快速生成专业的 Python 脚本和应用代码。一键创建完整项目结构,支持CLI、API、爬虫、Bot、Django等多种项目类型,包含完整的项目结构、配置文件、依赖管理、测试、README和文档。
async for message in pubsub.listen():
if message["type"] == "message":
print("收到:", message["data"])
- 必须在
await pubsub.subscribe("channel1")之后再调用listen() - 该迭代器不会自动退出,需配合
asyncio.create_task()启动为后台任务,或用asyncio.wait_for()控制超时 - 若中途想退订,调用
await pubsub.unsubscribe("channel1"),但当前正在 listen 的迭代器仍会收到 "unsubscribe" 类型消息
订阅多个频道时要注意消息类型判断
当同时 subscribe("a", "b"),listen() 返回的消息里 message["channel"] 是 bytes 类型(如 b"a"),且 message["type"] 不只是 "message",还可能有 "subscribe"、"unsubscribe"、"pmessage"(模式匹配时)等。
- 务必检查
message["type"] == "message"再处理业务逻辑,否则会把连接状态消息当业务数据解析出错 - 如果用
psubscribe("chan*"),则收消息时message["pattern"]和message["channel"]都存在,且message["type"]是"pmessage" - 所有字符串字段(
data,channel,pattern)默认是 bytes,需手动.decode("utf-8"),除非初始化时传decode_responses=True
生产环境必须处理连接中断与重连
aioredis 的 PubSub 对象不具备自动重连能力。网络抖动或 Redis 重启后,listen() 会抛出 ConnectionError 或直接静默断开,后续不再收消息。
- 不能依赖 try/except 包裹单次
async for循环——异常发生时迭代器已失效,必须重建pubsub实例并重新subscribe - 建议封装成带指数退避的重连循环,每次失败后
await asyncio.sleep(min(2**attempt, 60)) - 注意:重连期间发布的消息会丢失,Redis pubsub 本身不保证消息持久化
真正难的不是写通第一版,而是让 pubsub.listen() 在进程生命周期内稳定跑满数天不掉线——这需要对连接状态、异常类型、重试边界有足够细的控制。
Python免费学习笔记(深入):立即使用
在学习笔记中,你将探索 Python 的核心概念和高级技巧!










