Python 异步编程与 幂等设计 实践
前段时间把一个跑了两三年的消息回调服务用 asyncio 重构了一遍,压测数据确实好看,QPS 从三百多干到两千多,延迟也降了一截,当时还觉得自己挺牛。结果上线不到一周就翻车了:下游对账发现同一笔退款被处理了两次。翻日志才发现,是 MQ 那边网络抖动导致消费超时,消息被重新投递了一次,而我的消费者是"收到就干活",压根没管这条消息之前是不是已经处理过。
后来花了几天把幂等这块儿补上,顺便把异步场景下各种重复的来龙去脉捋了一遍。下面把做法和踩过的坑都记一下,代码都是实测能跑的。
异步环境里,重复根本躲不掉
一开始我天真地以为重复消费只是小概率事件,真去翻了文档才发现,主流的 MQ 基本都只保证 at least once,也就是消息至少送达一次,至于会不会多送,人家不管。再加上异步代码里还有几个天然的重复来源:
- 消费超时,broker 认为你挂了,消息重新入队
- 网络重试,客户端没收到 ack,自己又发了一遍
- 服务重启,
asyncio任务没跑完就被 kill,重启后从头再来 - 上游服务超时重试,webhook 场景特别常见
所以别指望"只要代码写得对就不会重复",正确的心态是:重复一定会发生,你的代码必须在重复发生时也能给出正确的结果。这就是幂等,说人话就是同一个操作执行一次和执行一百次,最终效果一样。网上不少文章吹 Kafka 能做到 exactly once,那也只是 Kafka 内部的事务语义,放到你业务的全链路里看,该重复还是会重复,幂等还是得自己做。
最稳的一层:数据库唯一索引
我实践下来最可靠的兜底就是数据库的唯一约束。给业务表加一个幂等键字段,重复插入直接报错,天然原子,不存在竞态:
CREATE TABLE refund_records (
id BIGSERIAL PRIMARY KEY,
idem_key VARCHAR(64) NOT NULL UNIQUE,
order_no VARCHAR(32) NOT NULL,
amount NUMERIC(10, 2) NOT NULL,
created_at TIMESTAMPTZ DEFAULT now()
);
消费端用 asyncpg 的话大概是这样:
async def handle_refund(conn, msg: dict):
try:
async with conn.transaction():
await conn.execute(
"INSERT INTO refund_records (idem_key, order_no, amount) VALUES ($1, $2, $3)",
msg["idem_key"], msg["order_no"], msg["amount"],
)
await conn.execute(
"UPDATE accounts SET balance = balance + $1 WHERE order_no = $2",
msg["amount"], msg["order_no"],
)
except asyncpg.UniqueViolationError:
# 这条消息已经处理过了,直接跳过
return
注意两个操作要放在同一个事务里,这样"插记录"和"改余额"要么都成功要么都不成功,不会出现记录插了但钱没动的尴尬局面。MQ 的 ack 放在事务提交之后,这个顺序别搞反了。
快的那一层:Redis SETNX 前置拦截
数据库唯一索引虽然稳,但每次重复请求都要打到数据库,量大的时候有点浪费。所以在前面再加一层 Redis 拦截,一条命令就能挡掉大部分重复:
from redis.asyncio import Redis
redis = Redis.from_url("redis://127.0.0.1:6379/0")
async def try_acquire(idem_key: str) -> bool:
return await redis.set(f"idem:{idem_key}", "1", nx=True, ex=86400)
关键是 SET key value NX EX ttl,设置成功说明是第一次来,返回 False 就是重复请求。千万别写成先 GET 判断再 SET,两步之间另一个协程可能已经插进去了,事件循环里这种竞态比多线程还隐蔽,因为你想不到代码会在哪个 await 处切走。我这里用的是 redis 这个库自带的 asyncio 支持,老项目还在用 aioredis 的话也没问题,用法基本一样,不过这库已经合并进官方了,新项目没必要再装。
一个完整的 webhook 例子
把上面两层组合起来,一个带幂等的 aiohttp 接口大概长这样:
async def refund_webhook(request: web.Request) -> web.Response:
idem_key = request.headers.get("X-Idempotency-Key")
if not idem_key:
return web.json_response({"error": "idempotency key required"}, status=400)
if not await try_acquire(idem_key):
return web.json_response({"msg": "duplicated request"})
try:
payload = await request.json()
async with request.app["pg"].acquire() as conn:
await handle_refund(conn, payload)
except Exception:
# 明确失败了才删幂等键,让上游可以重试
await redis.delete(f"idem:{idem_key}")
raise
return web.json_response({"msg": "ok"})
有个小细节:只有业务处理明确失败的时候才删幂等键。处理成功的话键就留着,等 TTL 到期自然过期,这段时间内的重复请求都会被挡掉。
踩过的几个坑
坑一:任务被 cancel 时乱删幂等键。 asyncio 里任务被取消是家常便饭,优雅关停、超时控制都会触发 CancelledError。我一开始在 finally 里无条件删键,结果任务跑到一半被取消,数据库事务其实已经提交了,键却被删了,下一条重复消息进来又执行了一遍退款。后来改成只在捕获到明确异常时才删,宁可让重复请求被多挡一会儿。顺便说一句,3.8 之后 CancelledError 继承自 BaseException,上面那段 except Exception 捕获不到它,这反而帮了忙。
坑二:幂等键的来源要想清楚。 最好让上游在请求头里带 X-Idempotency-Key,一般用 UUID。上游不给的话也可以自己对报文做 hash,但要小心:有些业务本来就允许内容相同的多次请求,比如用户今天买一杯奶茶、明天又买了杯一模一样的,拿内容 hash 当幂等键就把合法的第二次请求挡掉了,这种时候把用户 ID 和时间窗口拼进 hash 会好一点。
坑三:TTL 别拍脑袋定。 Redis 那层的过期时间要覆盖上游可能重试的最大窗口,对账相关的业务我一般 24 小时起步,太短了重复请求会穿透到数据库层,虽然兜得住,但没必要。
坑四:ack 的时机。 消费 MQ 时 ack 一定放在幂等写入或事务提交之后,反过来的话,业务刚执行完还没来得及 ack 就崩了,消息重投,虽然幂等键兜得住,但白白多走一遍流程。
最后
这一轮补完之后跑了一个多月,下游再没报过账对不上的问题。总结下来其实就一句话:异步代码里,把每一次操作都当成可能会被执行两次来写,数据库唯一索引做最后防线,Redis 挡在前面省点钱,剩下的就是耐心处理各种异常路径。说起来不复杂,但都是被重复扣款教育过才记得住的。
评论
还没有评论。
发表评论
提交后评论将经过自动审核,审核通过后公开展示。