推荐使用 redis.asyncio(aioredis 已合并入 redis-py v4+),需显式设 decode_responses=true、用 async with 管理连接、避免混用同步客户端;lpush+brpop 实现轻量 fifo 队列,高可靠场景应选 stream+消费组。

用 aioredis 建立异步 Redis 连接
同步 Redis 客户端(如 redis-py)在 await 表达式里会阻塞事件循环,必须换用原生支持 asyncio 的客户端。截至 2026 年,aioredis 已被官方合并进 redis-py v4+,所以直接安装并导入 redis.asyncio 即可:
pip install redis>=4.6.0
连接时务必用 async with 管理生命周期,避免连接泄漏:
import redis.asyncio as redis <p>async def get_client(): async with redis.Redis(host="localhost", port=6379, db=0) as client: await client.ping() # 验证连通性 return client</p>
- 别漏掉
db=0参数——不显式指定时,默认 db 是 0,但某些环境(如云 Redis)可能限制访问非 0 db,导致ConnectionError -
socket_timeout=3建议显式设置,否则网络抖动时协程会无限等待 - 不要复用全局 client 实例;
async with每次新建连接更安全,尤其在高并发短任务场景下
用 LPUSH + BRPOP 实现带阻塞的 FIFO 队列
Redis List 是最轻量、最可控的队列载体,LPUSH 入队、BRPOP 出队天然支持阻塞等待,避免轮询浪费 CPU。
注意:必须用 BRPOP 而不是 RPOP,否则消费者会忙等(busy-loop),且无法感知新消息到达:
async def producer(client: redis.Redis, queue_name: str, message: str):
await client.lpush(queue_name, message)
<p>async def consumer(client: redis.Redis, queue_name: str):</p><h1>timeout=0 表示永久阻塞,生产环境建议设为 5~30 秒</h1><pre class="brush:python;toolbar:false;">result = await client.brpop(queue_name, timeout=5)
if result is None:
return None # 超时,可做心跳或重试逻辑
_, payload = result # BRPOP 返回 (key, value),丢弃 key
return payload
-
BRPOP是**单 key** 命令,不支持同时监听多个队列;若需多队列分发,得用多个协程或改用 Stream - 入队用
LPUSH、出队用BRPOP才是严格 FIFO;反过来(RPUSH+BLPOP)也行,但混用会导致顺序错乱 - 消息体建议 JSON 序列化,避免二进制数据引发解码异常:
await client.lpush("q", json.dumps({...}).encode())
用 redis.Stream 替代 List 实现可靠消费组
当需要消息不丢失、支持多消费者协作、记录消费进度时,List 就力不从心了。Stream 是 Redis 5.0+ 提供的专为消息队列设计的数据结构,配合消费组(consumer group)能解决确认、重试、负载均衡问题。
Redis 缓存和数据结构管理技能。通过自然语言操作 Redis,支持 String、Hash、List、Set、ZSet、Stream 等数据结构操作。当用户提到 Redis、缓存、消息队列、会话存储时使用此技能。
关键命令对应关系:
- 生产者 →
XADD stream_name * field1 value1 - 消费者组首次读取 →
XREADGROUP GROUP mygroup consumer1 STREAMS mystream > - 消息处理完后确认 →
XACK mystream mygroup <id></id>
Python 中使用示例:
async def produce_to_stream(client, stream_name, data):
await client.xadd(stream_name, data)
<p>async def consume_from_group(client, stream_name, group_name, consumer_name):</p><h1>> 表示只读取新消息;也可以用 ID 读历史未确认消息</h1><pre class="brush:python;toolbar:false;">messages = await client.xreadgroup(
group_name, consumer_name,
streams={stream_name: ">"},
count=1,
block=5000 # 单位毫秒,比 BRPOP 的秒级更精细
)
if not messages:
return []
return messages[0][1] # 解包结构
- 首次创建消费组必须先调用
XGROUP CREATE,否则xreadgroup报NOGROUP错误 -
block参数单位是毫秒,不是秒;设为 0 会永久阻塞,不推荐 - Stream 消息 ID 是时间戳+序列号(如
1725578001234-0),不能手动构造,必须由 Redis 生成
避免在 asyncio 中混用同步 Redis 和线程池
有人试图用 loop.run_in_executor 包裹同步 redis.Redis 调用,这看似“异步”,实则埋雷:
- 每个 executor 线程都持有一个独立 TCP 连接,Redis 连接数会随并发线程数线性上涨,很快打满 maxclients
- 连接池无法跨线程共享,
ConnectionPool在线程内失效,失去复用意义 - 异常堆栈被 executor 吞掉,debug 困难;比如
TimeoutError可能被转成concurrent.futures.TimeoutError,掩盖真实来源
真正该做的,是坚持全链路异步:
- 所有 Redis 操作走
redis.asyncio.Redis - HTTP 请求用
aiohttp或httpx.AsyncClient,不用requests - 文件 I/O 用
asyncio.to_thread(open, ...)(Python 3.9+)或anyio,而非自建线程池
Stream 的 ACK 机制和消费组偏移管理,是 List 永远无法提供的能力;但如果你只需要“发完即忘”的日志收集或通知广播,List + BRPOP 就足够轻、足够快——选型前先问清楚:消息丢了能不能接受?要不要多个 worker 分摊压力?有没有重试需求?
Python免费学习笔记(深入):立即使用
在学习笔记中,你将探索 Python 的核心概念和高级技巧!










