支付系统的核心诉求只有四个字:不丢单、不错账。用户付款成功的瞬间,微信/支付宝会立刻发起异步回调,如果我们在回调接口里同步完成验签、改单、扣库存、发卡密、发通知,任何一个环节抖动都会拖垮整个下单链路;回调超时后平台还会反复重试,把订单状态搅成一团。本文分享用消息队列(MQ)重构支付链路的完整方案:回调异步化、可靠投递、延迟关单、最终一致性,全部附可运行的 Python 代码。

一、为什么支付回调必须异步化

先看一个典型的"同步噩梦"——所有逻辑都写在回调接口里:

# 同步版:回调里什么都干,慢且脆弱
def pay_callback(request):
    verify_sign(request)          # 1. 验签
    update_order(request)         # 2. 改订单状态
    deduct_stock()                # 3. 扣库存
    send_card_code()              # 4. 发卡密(可能调第三方短信)
    notify_user()                 # 5. 发通知
    return "success"

这段代码有三个致命问题:——发卡密要调第三方接口,动辄几百毫秒,回调接口整体变慢,支付平台会不断重试;重复——支付宝/微信回调最多重试 15 次,每次进来都重新扣一次库存;耦合——发货逻辑和支付逻辑写死在一起,换发货渠道就要改回调代码。用 MQ 改造后,回调接口只做两件事:验签 + 投递消息,立即返回 success,重活全部丢给消费者慢慢做。

二、选型:Redis Stream、RabbitMQ 还是 Kafka

小型支付项目我首推 Redis Stream,因为 Redis 本来就在用,零新增依赖;中等规模用 RabbitMQ,延迟队列开箱即用;海量数据才上 Kafka。对比如下:

维度Redis StreamRabbitMQKafka
部署成本低(复用 Redis)中(Erlang 环境)高(KRaft/ZooKeeper)
延迟队列需 ZSET 自己实现TTL + 死信交换机需插件/自研
消息回溯支持不支持支持(按 offset)
吞吐量万级/秒万级/秒百万级/秒
适合场景中小支付站业务消息、延迟任务大数据、日志、对账

三、场景一:回调异步化 + 转卡码自动发货

以转卡码系统为例:用户支付成功后,系统要给订单绑定一张卡密并展示给用户。改造后的回调接口只要 10 毫秒就能返回:

# 回调接口:只验签 + 投递,秒回 success
@app.post("/api/pay/callback")
def pay_callback():
    body = request.get_data()
    if not verify_alipay_sign(body):       # 验签失败直接拒绝
        return "fail", 400
    msg_id = hashlib.md5(body).hexdigest() # 用报文做幂等键
    r.xadd("stream:pay", {"msg_id": msg_id, "body": body.decode()})
    return "success"

消费者从消费组取消息,先查幂等表再执行发货,独立进程、可水平扩容:

# 消费者:消费组保证一条消息只被一个实例处理
while True:
    msgs = r.xreadgroup("g:pay", "c1", {"stream:pay": ">"}, count=10, block=5000)
    for _, items in msgs:
        for msg_id, data in items:
            try:
                if not order_done(data["msg_id"]):   # 幂等检查
                    bind_card_to_order(data)          # 绑定卡密
                    mark_order_paid(data["msg_id"])
                r.xack("stream:pay", "g:pay", msg_id) # 确认消费
            except Exception:
                r.xadd("stream:pay:retry", data)      # 进重试队列

这样回调接口永远秒回,平台不会再重试;消费者挂了消息也不丢(Stream 持久化 + PEL 待确认列表);想提升发货速度,多开几个消费者即可。卡密库存的扣减建议用 Redis 原子操作(DECR + 校验非负),避免并发超卖。

四、场景二:延迟队列实现订单超时自动关单

用户下单后 15 分钟未支付,订单要自动关闭并释放卡密库存。用 Redis ZSET 实现延迟队列只要 20 行:

# 下单时:把订单丢进延迟队列,score = 到期时间戳
r.zadd("delay:order", {order_id: time.time() + 900})
# 定时扫描线程:每秒取一次到期订单
def close_expired_orders():
    now = time.time()
    for oid in r.zrangebyscore("delay:order", 0, now, start=0, num=100):
        if r.zrem("delay:order", oid):   # 原子移除,防止重复处理
            if order_is_unpaid(oid):
                close_order(oid)
                release_stock(oid)       # 释放卡密库存

用 RabbitMQ 更简单:给队列设置 TTL 并绑定死信交换机,消息过期自动进入关单队列。注意关单前必须再次确认订单仍是"待支付"状态,防止用户刚好在最后一秒完成支付,出现"钱付了、单关了"的资损事故。关单后若收到迟到的支付回调,要走退款逆向流程。

五、场景三:失败重试与死信队列

发货失败(比如卡密库存不足)不能无限重试。标准做法是分级重试:前 5 次间隔递增,超过次数进死信队列人工介入:

# 重试计数器 + 死信判定
def consume(msg):
    try:
        deliver_card(msg)
    except Exception:
        retry = r.hincrby("retry:count", msg["msg_id"], 1)
        if retry <= 5:
            r.zadd("delay:retry", {msg["msg_id"]: time.time() + 60 * retry})
        else:
            r.xadd("stream:dead", msg)   # 死信队列,告警人工处理
            alert_admin(msg)

死信队列要接上告警(企业微信/钉钉机器人),人工处理后补发卡密,并在后台留操作日志,方便后续复盘。

六、最终一致性:本地消息表兜底

MQ 再可靠也可能丢消息(比如消费者消费后、确认前宕机)。生产环境我会加一道"本地消息表 + 定时任务"兜底:发货时在同一数据库事务里写订单状态和待发消息表,定时任务扫描未确认的消息重新投递;每日再跑对账脚本比对订单与支付流水。三层保障叠加,才能做到真正不丢单、不错账。

七、总结

消息队列不是银弹,但用在支付链路"回调 → 发货 → 通知"这一段,收益立竿见影:接口响应快、系统可水平扩容、故障可重试可追溯。如果你不想从零搭建,源码商城的转卡码系统 V3 已内置 Redis Stream 异步发货与延迟关单模块,开箱即用;开发调试时配合 Codex Desktop 这类 AI 编程助手,半小时就能把整套 MQ 链路跑起来。记住架构铁律:回调接口永远秒回,重活交给异步,兜底交给对账