Python 基础体系 · 第 86/112 篇。示例统一以 Python 3.14 为语言基线;第三方库使用与其兼容的现代稳定版本,版本敏感行为会单独说明。
Python 与 Redis:连接池、Pipeline、事务、缓存和分布式锁
Redis 是一个通过网络访问的内存数据服务。Python 应用并不是“调用一个本地字典”,而是把命令编码后发送给 Redis 服务器,再等待服务器返回结果。这个事实决定了本文的主线:
- 连接池解决连接生命周期和并发复用问题;
- Pipeline减少网络往返;
- 事务保证一组命令在 Redis 服务端连续执行;
- 缓存解决重复读取和后端压力问题;
- 分布式锁协调多个进程、线程或机器对共享资源的访问。
Python 侧通常使用 redis-py,安装包名为 redis。它既提供同步 API,也提供基于 asyncio 的异步 API。Redis 官方文档将连接池、Pipeline、事务、异步客户端和锁分别作为独立能力介绍,因此不能把它们简单视为同一个“性能优化开关”。(redis.io)
一、先建立正确的模型:Redis 客户端、连接和命令
一个最小的同步示例是:
import redis
r = redis.Redis(
host="localhost",
port=6379,
db=0,
decode_responses=True,
)
print(r.ping())
r.set("user:1:name", "Alice")
print(r.get("user:1:name"))
预期输出:
True
Alice
这里有三个对象层次:
Redis客户端对象:提供get()、set()、pipeline()等 Python 方法;- TCP 连接:客户端通过它与 Redis 服务器通信;
- Redis 命令:例如
GET user:1:name和SET user:1:name Alice。
decode_responses=True 表示把 Redis 返回的字节串解码为字符串。关闭它时,GET 通常返回 bytes,例如 b"Alice"。这不是数据是否存在的区别,而是客户端的响应解码策略。
Redis 中常见的数据类型包括:
| Redis 类型 | 典型用途 | Python 示例 |
|---|---|---|
| String | 字符串、数字、JSON 文本、锁令牌 | set()、get() |
| Hash | 对象字段、计数器集合 | hset()、hgetall() |
| List | 简单队列、栈 | lpush()、rpop() |
| Set | 去重集合、集合运算 | sadd()、smembers() |
| Sorted Set | 排行榜、延迟队列 | zadd()、zrange() |
| Stream | 消息流、消费组 | xadd()、xreadgroup() |
Redis 的 key 没有传统数据库中的表结构,因此应用必须自行设计命名空间。例如:
cache:user:42
lock:order:1001
rate-limit:ip:192.0.2.1
celery:task:<id>
key 命名不仅影响可读性,还影响批量扫描、监控、租户隔离和故障清理。生产环境中应避免使用过于宽泛的 KEYS *,因为它可能阻塞 Redis 处理其他请求;排查和遍历通常使用 SCAN。
二、连接池:复用 TCP 连接,而不是复用一次请求
2.1 为什么需要连接池
建立 TCP 连接涉及网络握手、认证以及客户端和服务器状态初始化。假设每个请求都执行:
def get_user(user_id: int) -> str | None:
r = redis.Redis(host="localhost", port=6379)
try:
return r.get(f"user:{user_id}:name")
finally:
r.close()
这段代码虽然可能“能工作”,但每次调用都创建客户端并触发连接管理。高并发下,连接建立和释放会产生额外开销,还可能导致 Redis 的连接数快速增长。
连接池是一组可复用的 Redis 连接。客户端从池中取得连接,命令完成后把连接归还,而不是立即销毁。Redis 官方建议生产应用使用连接池;异步客户端文档进一步建议在长生命周期应用启动时创建一个客户端,并在请求和任务之间共享它,而不是每个请求创建新的 Redis()。(redis.io)
同步客户端可以这样创建:
import redis
pool = redis.ConnectionPool.from_url(
"redis://localhost:6379/0",
max_connections=20,
decode_responses=True,
)
r = redis.Redis(connection_pool=pool)
print(r.ping())
这里:
max_connections=20是该进程最多使用的连接数;decode_responses=True控制响应解码;Redis(connection_pool=pool)创建使用这个池的客户端。
一个重要边界是:连接池限制的是单个 Python 进程中的 Redis 连接数。如果部署了 8 个进程,每个进程的 max_connections=20,理论上总连接上限接近:
其中:
- 是进程数;
- 是每个进程的连接上限。
因此示例中的总连接数约为:
这还没有计入 Celery Worker、定时任务进程、管理脚本和其他服务。
2.2 连接池大小不是越大越好
如果 Redis 每个命令平均耗时为 ,单连接在理想情况下每秒最多完成约:
但应用吞吐还受到 Redis 单线程命令执行、网络延迟、序列化和后端处理能力的限制。无限增大连接池不能无限提升吞吐,反而会增加:
- Redis 连接管理开销;
- 服务端连接数;
- 应用上下文切换;
- 突发时对 Redis 的压力;
- 连接超时和排队难以诊断的问题。
连接池大小应由并发模型决定。一个同步 Web Worker 通常在阻塞等待 Redis 时占用线程;异步应用可以让其他协程运行,但每条正在执行或等待响应的 Redis 操作仍需要合适的连接资源。
2.3 客户端和连接池的生命周期
在 FastAPI 等 ASGI 应用中,Redis 客户端适合放在应用生命周期中创建和关闭,而不是放到请求函数内部。FastAPI 推荐使用 lifespan 参数管理启动和关闭逻辑;ASGI Lifespan 协议要求资源初始化发生在接收请求前,清理发生在应用关闭时,并且每个处理请求的事件循环拥有自己的生命周期。(fastapi.tiangolo.com)
下面是一个可运行的异步 FastAPI 示例:
from contextlib import asynccontextmanager
from typing import AsyncIterator
import redis.asyncio as redis
from fastapi import FastAPI, Request
@asynccontextmanager
async def lifespan(app: FastAPI) -> AsyncIterator[None]:
client = redis.Redis.from_url(
"redis://localhost:6379/0",
decode_responses=True,
max_connections=20,
socket_connect_timeout=1.0,
socket_timeout=1.0,
)
try:
await client.ping()
app.state.redis = client
yield
finally:
await client.aclose()
app = FastAPI(lifespan=lifespan)
@app.get("/health/redis")
async def redis_health(request: Request) -> dict[str, bool]:
client: redis.Redis = request.app.state.redis
return {"redis": bool(await client.ping())}
运行前需要:
python -m pip install fastapi redis uvicorn
uvicorn app:app --reload
请求:
curl http://127.0.0.1:8000/health/redis
预期输出:
{"redis":true}
这里的错误路径也很重要:如果 await client.ping() 在启动阶段失败,应用可以直接启动失败,而不是启动后所有请求都返回 Redis 连接错误。finally 保证关闭阶段尝试释放客户端资源。
在多进程部署中,每个进程都会执行自己的 lifespan,因此每个进程应创建自己的异步 Redis 客户端,不能把一个绑定到事件循环的异步连接池跨进程或跨事件循环共享。ASGI 规范明确说明,多进程环境中每个进程都有自己的 lifespan。(asgi.readthedocs.io)
三、Pipeline:把多条命令批量发送
3.1 Pipeline 解决什么问题
假设应用依次执行 100 次 SET:
for i in range(100):
r.set(f"item:{i}", i)
如果每次调用都产生一次请求—响应往返,那么网络往返次数大致是:
Pipeline 将命令先缓存在客户端,调用 execute() 时批量发送:
with r.pipeline(transaction=False) as pipe:
for i in range(100):
pipe.set(f"item:{i}", i)
results = pipe.execute()
print(len(results))
print(results[:3])
预期输出类似:
100
[True, True, True]
Redis 官方定义中,Pipeline 的主要收益是减少网络和处理开销:多条命令一次发送,服务器一次返回结果。命令只有在 execute() 时才真正执行,Pipeline 中的命令方法返回的是 Pipeline 对象,因此可以链式调用。(redis.io)
Pipeline 主要改变的是通信形态:
普通模式:
客户端 -> SET -> Redis
客户端 <- OK <- Redis
客户端 -> SET -> Redis
客户端 <- OK <- Redis
Pipeline:
客户端 -> SET, SET, SET, SET -> Redis
客户端 <- OK, OK, OK, OK <- Redis
它不等于“把所有命令变成一个业务事务”。是否具有事务语义,取决于是否启用 transaction。
3.2 Pipeline 的内存边界
命令会在客户端缓冲。如果一次 Pipeline 放入数百万条命令,客户端需要保存命令参数和结果,可能造成内存压力。更稳妥的方式是分块:
def batch_set(client: redis.Redis, values: dict[str, str], batch_size: int = 500):
items = list(values.items())
for start in range(0, len(items), batch_size):
batch = items[start:start + batch_size]
with client.pipeline(transaction=False) as pipe:
for key, value in batch:
pipe.set(key, value)
pipe.execute()
分块之后:
其中 是每批命令数,而不是所有数据量 。代价是网络往返次数变为:
所以批大小是在客户端内存、单次请求大小和网络往返次数之间做取舍。
四、事务:连续执行,不等于数据库式回滚
4.1 Redis 事务的实际语义
Redis 事务通常由以下命令组成:
MULTI
command-1
command-2
EXEC
事务的核心保证是:EXEC 执行后,事务中的命令会连续执行,不会被其他客户端命令插入。Redis 不提供传统关系数据库事务常见的“任意一步失败就自动回滚”语义。
在 redis-py 中:
with r.pipeline(transaction=True) as pipe:
pipe.set("account:1:balance", 900)
pipe.set("account:2:balance", 1100)
results = pipe.execute()
print(results)
redis-py 的 Pipeline 默认启用事务;如果只想批量发送而不使用事务,可以显式设置 transaction=False。(redis.io)
4.2 “原子执行”到底是什么意思
考虑两个客户端 A、B:
A: MULTI
A: SET x 1
A: SET y 1
A: EXEC
B: SET x 9
如果 A 的 EXEC 先开始执行,Redis 会连续执行 A 的两条命令,然后才处理 B 的命令。B 不会插入到 A 的两条命令之间。
但是,下面这些情况不会自动回滚:
MULTI
SET user:name Alice
HSET user:name age 20
EXEC
如果 user:name 是字符串,第二条 HSET 会产生类型错误。第一条 SET 已经执行成功,Redis 不会撤销它。事务结果可能是:
[OK, WRONGTYPE ...]
因此,Redis 事务的正确理解是:
它保证命令序列在服务端连续执行,但不保证失败后的自动回滚。
4.3 WATCH:基于版本变化的乐观锁
如果事务中的新值依赖于事务开始前读到的数据,就需要 WATCH。
例如实现“余额增加”:
from redis.exceptions import WatchError
def add_balance(client: redis.Redis, key: str, amount: int) -> int:
while True:
try:
with client.pipeline() as pipe:
pipe.watch(key)
raw = pipe.get(key)
current = int(raw or 0)
new_value = current + amount
pipe.multi()
pipe.set(key, new_value)
result = pipe.execute()
return new_value
except WatchError:
# key 在 WATCH 之后被其他客户端修改,重新读取并计算
continue
它的状态变化如下:
初始:balance = 100
客户端 A WATCH balance
客户端 A GET balance -> 100
客户端 B SET balance 150
客户端 A MULTI
客户端 A SET balance 110
客户端 A EXEC -> WatchError
A 不能把旧值 100 加上 10 后覆盖 B 的 150,因为 A 观察到 key 已经变化。A 必须重新读取:
重新读取:balance = 150
重新计算:150 + 10 = 160
重新执行:SET balance 160
WATCH 的形式化条件是:
也就是:事务执行时,key 的版本必须仍然等于读取时的版本。如果不成立,事务失败并重试。
需要注意三点:
WATCH失败不是 Redis 服务器故障,而是并发冲突;- 高竞争 key 可能导致大量重试;
- 重试必须有上限、退避和监控,否则会形成忙等。
如果更新逻辑可以由 Redis 原生命令表达,优先使用原子命令。例如计数器不需要:
current = int(r.get("counter") or 0)
r.set("counter", current + 1)
因为两个客户端可能同时读到 10,最后都写入 11。应使用:
r.incr("counter")
如果逻辑无法由单条命令表达,可以考虑 WATCH 或 Lua 脚本。Lua 脚本把读取、判断和写入放进 Redis 服务端连续执行,避免 Python 代码在网络往返期间暴露并发窗口。
五、Pipeline、事务、Lua 和锁的区别
可以用四个问题区分它们:
| 机制 | 主要解决的问题 | 是否减少往返 | 是否保证连续执行 | 是否自动回滚 |
|---|---|---|---|---|
| Pipeline | 批量发送命令 | 是 | 取决于 transaction |
否 |
| Redis 事务 | 一组命令连续执行 | 通常是 | 是 | 否 |
WATCH |
检测读改写冲突 | 不一定 | 成功时是 | 否,失败重试 |
| Lua | 在 Redis 服务端执行复合逻辑 | 是 | 是 | 否 |
| 分布式锁 | 协调多个客户端的临界区 | 不以此为目标 | 取决于实现 | 不适用 |
例如,下面这段代码使用 Pipeline,但没有事务保证:
with r.pipeline(transaction=False) as pipe:
pipe.set("a", 1)
pipe.set("b", 2)
pipe.execute()
它减少了网络往返,但其他客户端可以在 a 和 b 的执行之间修改数据。
而下面这段使用默认事务 Pipeline:
with r.pipeline(transaction=True) as pipe:
pipe.set("a", 1)
pipe.set("b", 2)
pipe.execute()
它既批量发送,又使用 MULTI/EXEC,因此两条命令在服务端连续执行,但依旧没有失败回滚。
六、缓存:Redis 不是“把数据库复制一份”
6.1 缓存的基本模型
缓存是一个速度更快、通常不作为最终事实来源的数据副本。最常见的是 Cache-Aside,也称旁路缓存:
flowchart TD
A[客户端请求] --> B[应用查询 Redis]
B -->|命中| C[返回缓存数据]
B -->|未命中| D[查询主数据库]
D --> E[写入 Redis 并设置 TTL]
E --> F[返回数据库数据]
G[业务写请求] --> H[更新主数据库]
H --> I[删除或更新缓存]
读流程:
import json
from typing import Callable, Any
def get_user(
client: redis.Redis,
user_id: int,
load_from_db: Callable[[int], dict[str, Any] | None],
) -> dict[str, Any] | None:
key = f"cache:user:{user_id}"
cached = client.get(key)
if cached is not None:
return json.loads(cached)
user = load_from_db(user_id)
if user is None:
return None
client.set(
key,
json.dumps(user, ensure_ascii=False),
ex=60,
)
return user
写流程:
def update_user(
client: redis.Redis,
user_id: int,
payload: dict[str, Any],
update_db: Callable[[int, dict[str, Any]], None],
) -> None:
update_db(user_id, payload)
# 数据库成功后删除缓存
client.delete(f"cache:user:{user_id}")
Redis 官方将 Cache-Aside 描述为:先查 Redis,未命中时查询主数据源,再将结果以 TTL 写回;更新主数据后删除缓存,从而限制过期数据窗口。(redis.io)
6.2 TTL 和缓存一致性
TTL,即 Time To Live,表示 key 还能存活多久。设置:
r.set("cache:article:1", '{"title":"Redis"}', ex=60)
print(r.ttl("cache:article:1"))
ex=60 表示 60 秒后过期。TTL 只能保证:
它不能保证数据在数据库更新后立即变新。假设缓存刚写入,TTL 为 60 秒,数据库在第 1 秒更新,那么缓存仍可能继续返回旧值约 59 秒。
因此:
- TTL 是过期上界,不是强一致保证;
- 对必须立即生效的数据,写数据库后应删除或更新缓存;
- 删除缓存也不是绝对无竞态,复杂写入场景要设计写入顺序和版本号。
一个常见竞态如下:
时刻 1:请求 A 从数据库读到旧值 old
时刻 2:请求 B 更新数据库为 new,并删除缓存
时刻 3:请求 A 把 old 写回缓存
此时缓存又被旧值覆盖。解决方式包括:
- 让缓存写入携带版本号,只允许新版本覆盖旧版本;
- 使用消息队列串行化更新;
- 使用数据库变更事件驱动失效;
- 对关键场景使用带版本判断的 Lua 脚本。
6.3 缓存击穿、穿透和雪崩
缓存击穿:某个热点 key 过期,大量请求同时查询数据库。
设瞬时并发为 ,单个数据库查询成本为 。如果没有保护,数据库瞬间承受约 次查询;如果只有一个请求负责回源,其余请求等待或读取新缓存,数据库压力接近 1 次回源。
缓存穿透:请求查询一个数据库中本来不存在的 ID,结果每次都无法写入有效缓存。可以缓存短 TTL 的空值:
NOT_FOUND = "__NOT_FOUND__"
cached = r.get(key)
if cached == NOT_FOUND:
return None
if cached is None:
user = load_from_db(user_id)
if user is None:
r.set(key, NOT_FOUND, ex=10)
return None
r.set(key, json.dumps(user), ex=60)
return user
return json.loads(cached)
空值 TTL 要短于正常数据,否则新创建的数据可能被旧的“不存在”缓存遮挡。
缓存雪崩:大量 key 在同一时间过期,导致集中回源。可以给 TTL 加随机抖动:
import random
ttl = 60 + random.randint(0, 15)
r.set(key, value, ex=ttl)
这只能分散过期时间,不能替代容量保护和回源限流。
6.4 用锁保护热点回源
对于热点 key,可以使用一个短时锁让一个请求回源:
import time
import uuid
def load_with_single_flight(
client: redis.Redis,
cache_key: str,
load_from_db,
) -> str | None:
lock_key = f"{cache_key}:fill-lock"
token = uuid.uuid4().hex
cached = client.get(cache_key)
if cached is not None:
return cached
acquired = client.set(lock_key, token, nx=True, ex=5)
if acquired:
try:
cached = client.get(cache_key)
if cached is not None:
return cached
value = load_from_db()
if value is not None:
client.set(cache_key, value, ex=60)
return value
finally:
release_lock(client, lock_key, token)
else:
# 简化示例:短暂等待后重读缓存
for _ in range(10):
time.sleep(0.05)
cached = client.get(cache_key)
if cached is not None:
return cached
return load_from_db()
这里的锁只保护“缓存填充”这一小段临界区,不应把整个长时间业务流程都放进去。真实系统还需要限制等待者数量,并处理回源失败、锁过期和 Redis 不可用等情况。
七、分布式锁:从“互斥”到“租约”
7.1 分布式锁的定义
进程内的 threading.Lock 只能协调同一进程中的线程。分布式锁通过共享服务协调多个进程或机器:
应用实例 A ─┐
应用实例 B ─┼── Redis ── lock:resource
Celery Worker ─┘
一个锁至少要表达:
- 锁是否已被占用;
- 谁持有锁;
- 锁何时自动失效;
- 释放操作不能误删其他持有者的锁。
7.2 错误实现:直接 DEL
错误示例:
if r.set("lock:report", "A", nx=True, ex=10):
try:
do_work()
finally:
r.delete("lock:report")
反例时序:
t0:A 获得锁,TTL=10 秒
t1:A 执行时间过长,锁过期
t2:B 获得同一个锁
t3:A finally 执行 DEL
t4:B 的锁被 A 删除
A 已经不是当前持有者,却删除了 B 的锁。
Redis 官方给出的基本改进是:使用不可预测的随机 token 作为锁值,并只在 key 当前值仍等于该 token 时删除。SET key value NX EX seconds 可以把“仅当不存在时设置”和“设置过期时间”放在一个命令中。(redis.io)
释放锁需要 Lua 脚本:
import uuid
import redis
RELEASE_LOCK = """
if redis.call('get', KEYS[1]) == ARGV[1] then
return redis.call('del', KEYS[1])
else
return 0
end
"""
def acquire_lock(
client: redis.Redis,
key: str,
ttl: int,
) -> str | None:
token = uuid.uuid4().hex
if client.set(key, token, nx=True, ex=ttl):
return token
return None
def release_lock(
client: redis.Redis,
key: str,
token: str,
) -> bool:
result = client.eval(RELEASE_LOCK, 1, key, token)
return result == 1
使用:
lock_key = "lock:report"
token = acquire_lock(r, lock_key, ttl=30)
if token is None:
print("未获得锁")
else:
try:
do_work()
finally:
release_lock(r, lock_key, token)
Lua 脚本的重要性在于“比较”和“删除”必须是一个不可插入的操作。如果拆成:
if r.get(key) == token:
r.delete(key)
那么在 GET 和 DEL 之间,锁可能过期并被另一个客户端重新取得,旧持有者仍可能删除新锁。
7.3 锁的 TTL 是租约,不是执行时间预测
设:
- :锁 TTL;
- :业务实际执行时间;
- :网络和调度延迟;
- :重试、暂停或垃圾回收造成的额外时间。
为了让单实例锁在业务期间持续有效,需要近似满足:
但 往往不是固定值。如果把 TTL 设置得过短,锁可能在业务完成前过期;设置过长,则持有者宕机后,其他请求需要等待更长时间。
可选方案:
- 给锁设置合理 TTL;
- 在确认仍持有 token 的前提下续租;
- 把临界区缩短;
- 让业务操作本身具备幂等性;
- 不把分布式锁当作数据库提交保证。
最重要的边界是:锁只能阻止遵守该锁协议的参与者同时进入临界区。如果另一个组件直接修改数据库、绕过锁,Redis 锁不会保护它。
7.4 redis-py 的 Lock API
redis-py 还提供了 Lock 抽象,可以配置 timeout、blocking 和 blocking_timeout 等参数;其文档说明该锁可用于跨进程、跨机器协调。(redis.readthedocs.io)
示例:
lock = r.lock(
"lock:invoice:1001",
timeout=30,
blocking_timeout=2,
)
if not lock.acquire():
raise TimeoutError("无法在 2 秒内获得锁")
try:
generate_invoice()
finally:
lock.release()
这类 API 能减少手写协议的错误,但调用者仍必须理解:
timeout到期后锁可能自动释放;- 业务执行超过 TTL 时不能假定自己仍持有锁;
release()失败应记录日志并告警;- Redis 故障时,锁状态和业务状态可能无法同时判断。
7.5 单实例锁和 Redlock
单 Redis 实例锁的安全性依赖该实例及其故障转移语义。如果 Redis 主节点在复制完成前宕机,锁状态可能在新主节点上不存在,另一个客户端可能重新取得锁。
Redlock 是一种在多个独立 Redis 实例上尝试获取锁的算法。Redis 官方文档将它描述为比单实例方案更复杂、面向多实例的分布式锁模式,同时也指出获取多数锁失败时应尽快释放已部分获得的锁。(redis.io)
不过,Redlock 不是所有业务的默认答案。对于支付、库存、订单状态等高价值操作,仅靠 Redis 锁通常不足以形成最终正确性保证。更可靠的设计通常还需要:
- 数据库唯一约束;
- 条件更新,例如
UPDATE ... WHERE version = ?; - 幂等键;
- 状态机;
- 操作记录和补偿任务。
锁解决“同时进入”的问题,幂等解决“重复执行”的问题,持久化约束解决“最终不能违反”的问题。三者不能相互替代。
八、一个完整的 FastAPI 缓存示例
下面将连接池、应用生命周期和 Cache-Aside 组合起来。示例用内存字典模拟数据库,真实环境中应替换为数据库访问函数。
from contextlib import asynccontextmanager
from typing import AsyncIterator
import json
import redis.asyncio as redis
from fastapi import FastAPI, HTTPException, Request
DATABASE = {
1: {"id": 1, "name": "Alice"},
}
@asynccontextmanager
async def lifespan(app: FastAPI) -> AsyncIterator[None]:
client = redis.Redis.from_url(
"redis://localhost:6379/0",
decode_responses=True,
max_connections=20,
socket_connect_timeout=1,
socket_timeout=1,
)
await client.ping()
app.state.redis = client
try:
yield
finally:
await client.aclose()
app = FastAPI(lifespan=lifespan)
def cache_key(user_id: int) -> str:
return f"cache:user:{user_id}"
@app.get("/users/{user_id}")
async def read_user(user_id: int, request: Request):
client: redis.Redis = request.app.state.redis
key = cache_key(user_id)
cached = await client.get(key)
if cached is not None:
return {
"source": "cache",
"user": json.loads(cached),
}
user = DATABASE.get(user_id)
if user is None:
# 对不存在结果设置短 TTL,避免缓存穿透
await client.set(key, "__NOT_FOUND__", ex=10)
raise HTTPException(status_code=404, detail="user not found")
await client.set(
key,
json.dumps(user),
ex=60,
)
return {
"source": "database",
"user": user,
}
@app.put("/users/{user_id}")
async def update_user(
user_id: int,
payload: dict[str, str],
request: Request,
):
client: redis.Redis = request.app.state.redis
if user_id not in DATABASE:
raise HTTPException(status_code=404, detail="user not found")
DATABASE[user_id].update(payload)
# 主数据更新成功后使缓存失效
await client.delete(cache_key(user_id))
return DATABASE[user_id]
首次请求 /users/1 的数据流是:
Redis GET -> 空
读取 DATABASE
Redis SET,TTL=60
返回 source=database
第二次请求:
Redis GET -> JSON
不访问 DATABASE
返回 source=cache
更新请求:
修改 DATABASE
DEL cache:user:1
下一次读取重新从数据库加载。这个示例刻意没有把缓存写入和数据库更新伪装成一个 Redis 事务,因为 Redis 事务不能覆盖外部数据库。数据库和 Redis 之间的双写一致性必须通过顺序、重试、消息或补偿机制设计。
九、失败路径、诊断和观测
Redis 相关故障不能只看“接口返回 500”。至少要区分以下几类:
9.1 连接失败
表现可能包括:
ConnectionError;TimeoutError;- 启动阶段
ping()失败; - 请求延迟突然升高。
诊断重点:
try:
await client.ping()
except Exception:
logger.exception("redis health check failed")
raise
应记录:
- Redis 地址和逻辑用途;
- 命令类型;
- 超时时间;
- 请求 trace ID;
- 是否为缓存读、缓存写、锁获取还是锁释放;
- 异常类型。
不要把密码、完整连接 URI 或锁 token 写入日志。
9.2 Pipeline 部分结果错误
Pipeline 返回的是按命令顺序排列的结果:
with r.pipeline(transaction=False) as pipe:
pipe.set("a", "1")
pipe.hset("a", mapping={"field": "value"})
results = pipe.execute()
第二条命令可能因为 key 类型错误而失败,但第一条命令已经成功。排查时必须逐项对应命令序号,不能把 Pipeline 当作“要么全成功、要么全失败”。
9.3 WATCH 重试过多
如果一个 key 被频繁修改,可能出现:
WatchError -> 重试
WatchError -> 重试
WatchError -> 重试
这说明并发冲突率高,而不一定说明 Redis 性能差。应记录:
- 业务 key;
- 重试次数;
- 最终是否成功;
- 单次事务耗时;
- 竞争来源。
重试循环必须设置上限:
from redis.exceptions import WatchError
for attempt in range(5):
try:
...
break
except WatchError:
if attempt == 4:
raise
9.4 缓存命中率不能单独证明缓存健康
应同时观察:
还要看:
- Redis 命令延迟;
- 回源请求数量;
- 热点 key;
- key 的 TTL 分布;
- 缓存填充失败数;
- 空值缓存命中数;
- 锁获取成功率;
- 锁等待时间;
- 数据库 P95/P99。
命中率高但缓存返回大量旧数据,仍然是业务错误;命中率低但数据具有高随机性,也可能是合理现象。
十、与 Celery 和幂等设计的关系
Redis 常被 Celery 用作 Broker 或结果后端,但这不意味着“Redis 中的任务消息”和“业务缓存”可以随意混用。
建议至少区分命名空间:
cache:user:42
lock:order:1001
celery-task-meta:<task-id>
Celery 任务可能因为网络断开、Worker 崩溃、ACK 时机或重试策略而再次执行。因此任务中的 Redis 操作不能只依赖“任务通常执行一次”。
例如,错误的任务写法:
def charge_order(order_id: int):
payment_api.charge(order_id)
r.set(f"order:{order_id}:charged", "1")
如果支付接口成功,但 Worker 在写 Redis 前崩溃,任务重试可能再次扣款。Redis 锁也不能单独解决这个问题,因为锁可能过期、进程可能暂停,外部支付系统也可能不受该锁约束。
更可靠的方式是让支付请求带业务幂等键:
idempotency-key = order:1001:charge
并在业务数据库中记录状态:
PENDING -> CHARGING -> SUCCEEDED
-> FAILED
Redis 可以加速状态查询、抑制短时间重复提交或协调非关键临界区,但最终状态和不可重复副作用应由具备持久化语义的系统保证。
十一、选择依据
可以按问题本身选择机制:
只需要执行一条命令
直接调用:
r.incr("counter")
优先使用 Redis 提供的原子命令,而不是在 Python 中手动读改写。
需要批量执行,但命令之间不要求事务语义
使用:
r.pipeline(transaction=False)
适合批量写入、批量读取和减少网络往返。
需要多条命令连续执行
使用默认事务 Pipeline:
r.pipeline(transaction=True)
但要记住:连续执行不等于失败回滚。
需要“读取后根据旧值计算新值”
使用 WATCH 重试、Lua 脚本,或重新设计为 Redis 原子命令。
需要给重复读取加速
使用 Cache-Aside,并明确:
- 数据源是谁;
- TTL 多长;
- 更新时如何失效;
- 空值是否缓存;
- 热点 key 如何防击穿;
- Redis 故障时是否降级。
需要跨进程协调临界区
使用带随机 token、TTL 和安全释放的分布式锁;临界区必须短,业务必须考虑锁过期和进程崩溃。
需要保证最终业务正确
不要只依赖缓存、Pipeline 或锁。应使用数据库约束、版本条件、幂等键和可恢复状态机,把 Redis 放在它擅长的“快速访问与协调”位置,而不是把它误当成整个业务一致性系统。
系列导航与关联阅读
- 系列入口:Python 完整学习路线:从语言模型、并发到 Web、数据、AI 与生产交付
- 上一篇:Python 异步数据库访问:连接池、事务、取消、并发和一致性
- 下一篇:Celery 任务队列:Broker、Worker、ACK、重试、定时和幂等
- 延伸:Python 可观测性:日志、指标、Trace、Context 和故障定位
官方资料
本文依据 Python 官方文档、相关 PEP 与生态项目官方文档重新梳理;正文、示例与工程清单由 WR BLOG 编写。

评论
0 条讨论