小Cの已经记不起来的博客

Python 异步编程与 幂等设计 实践

前段时间把一个跑了两三年的消息回调服务用 asyncio 重构了一遍,压测数据确实好看,QPS 从三百多干到两千多,延迟也降了一截,当时还觉得自己挺牛。结果上线不到一周就翻车了:下游对账发现同一笔退款被处理了两次。翻日志才发现,是 MQ 那边网络抖动导致消费超时,消息被重新投递了一次,而我的消费者是"收到就干活",压根没管这条消息之前是不是已经处理过。

后来花了几天把幂等这块儿补上,顺便把异步场景下各种重复的来龙去脉捋了一遍。下面把做法和踩过的坑都记一下,代码都是实测能跑的。

异步环境里,重复根本躲不掉

一开始我天真地以为重复消费只是小概率事件,真去翻了文档才发现,主流的 MQ 基本都只保证 at least once,也就是消息至少送达一次,至于会不会多送,人家不管。再加上异步代码里还有几个天然的重复来源:

所以别指望"只要代码写得对就不会重复",正确的心态是:重复一定会发生,你的代码必须在重复发生时也能给出正确的结果。这就是幂等,说人话就是同一个操作执行一次和执行一百次,最终效果一样。网上不少文章吹 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 挡在前面省点钱,剩下的就是耐心处理各种异常路径。说起来不复杂,但都是被重复扣款教育过才记得住的。

评论

还没有评论。

发表评论

提交后评论将经过自动审核,审核通过后公开展示。

未在播放