async for 不能直接用于 redis.pubsub(),因 pubsub 对象非异步可迭代;须用 async with 创建实例、await subscribe 后,再 async for msg in pubsub.listen() 获取消息,且需显式处理连接关闭、重连与超时。

async for 不能直接用在 redis.pubsub() 上
很多人一上来就写 async for msg in redis.pubsub(),结果报 TypeError: 'PubSub' object is not async iterable。因为 aioredis.PubSub 实例本身不支持 async for,它不是异步迭代器——你得手动调用 get_message() 或用 listen() 启动长连接循环。
正确启动订阅:必须用 listen() + await 驱动事件循环
listen() 返回一个异步生成器,这才是真正的入口。它会持续接收消息,直到连接断开或被取消。别用同步风格的 while True + time.sleep(),那会阻塞整个协程。
- 必须在
async with中创建 PubSub 实例(redis.pubsub()),否则连接资源不会自动释放 -
await pubsub.subscribe("channel1", "channel2")必须成功执行后才能listen() - 每条消息是字典,含
"type"(如"message"、"subscribe")、"channel"、"data"字段 -
"data"默认是bytes,若初始化时没设decode_responses=True,就得手动.decode()
示例片段:
Python 3.14.2是Python编程语言在2025年12月5日发布的稳定版本,属于3.14系列的第二个维护更新。该版本包含了18项修复,重点解决了多进程、数据类及正则表达式等模块的回归问题,并修复了CVE-2025-12084等安全漏洞。此版本标志着自由线程模式(移除GIL)正式获得官方支持,是Python发展的重要里程碑。
async def handle_pubsub():
async with redis.pubsub() as pubsub:
await pubsub.subscribe("notifications")
async for msg in pubsub.listen(): # ✅ 正确
if msg["type"] == "message":
print("Received:", msg["data"].decode())
如何安全退出监听并避免连接残留
Ctrl+C 或服务关闭时,async for 可能卡在 listen() 的下一次 await,导致连接没关、任务没 cancel。关键点是显式控制生命周期:
- 不要依赖
try/except KeyboardInterrupt粗暴中断——listen()内部可能正阻塞在 socket.recv() - 用
asyncio.create_task()启动监听,并在 shutdown 时task.cancel(),配合await task等待清理完成 -
pubsub.close()必须被 await(即await pubsub.close()),否则底层连接可能滞留 - 如果用了
ConnectionPool,记得全局池也要在应用退出时调await pool.disconnect()
生产环境必须处理的三个断裂点
Redis 订阅连接比普通 GET/SET 更脆弱,网络抖动、Redis 重启、超时断连都会让 listen() 报错或静默停止。它不会自动重连,也不会抛出明确异常来提醒你。
-
await pubsub.get_message(ignore_subscribe_messages=True)可用于非阻塞轮询,适合做健康检查,但频繁调用会增压 - 捕获
ConnectionError和asyncio.TimeoutError后,应重建PubSub实例并重新subscribe - 不要假设
listen()是“永远运行”的——加超时包装(如asyncio.wait_for(..., timeout=30))并重试,才是可靠做法
最易被忽略的是:PubSub 连接和普通 Redis 连接共用一个连接池,但它的空闲检测逻辑不同;默认 health_check_interval 对 PubSub 无效,必须单独兜底。
Python免费学习笔记(深入):立即使用
在学习笔记中,你将探索 Python 的核心概念和高级技巧!










