分布式事务解决方案:最佳实践(一个试用期后端仔的深夜踩坑实录)

MQ堵车了
2025-12-17 13:03
阅读 3951

上周五晚上 11 点半,我正窝在沙发上敲着键盘,咖啡杯里泡着第三包速溶咖啡——别笑,远程办公的好处就是不用穿裤子写代码。突然手机一震,企业微信弹出消息:“小张,明天上线前务必搞定订单-库存一致性的问题,不然双11大促怕是要翻车。”

我盯着屏幕,心里一万只草泥马奔腾而过。入职这家公司刚满三周,还在试用期,结果就被扔进了一个典型的“分布式事务”火坑。产品说要保证用户下单成功时库存一定扣减,否则就超卖;测试说他们压测时发现库存偶尔会多扣;运维大哥则在群里幽幽地说:“这要是线上崩了,锅可不小啊。”

于是,我这个 Python 后端新人,在深夜的台灯下,开启了一段与分布式事务斗智斗勇的旅程。


起因:微服务拆分后的“甜蜜烦恼”

我们公司的电商业务最近做了微服务重构:订单服务、库存服务、用户服务各自独立部署,通过 HTTP API 或消息队列通信。架构图看起来很美,但问题也随之而来。

比如用户下单这个操作,理想流程是:

  1. 订单服务创建订单(状态为“待支付”)
  2. 库存服务扣减商品库存
  3. 支付完成后,订单状态更新为“已支付”

但现实是:第 1 步成功了,第 2 步因为库存服务宕机失败了,结果用户看到订单创建成功,却无法支付——因为库存没扣,系统认为“无货”。更糟的是,如果第 2 步执行了但第 3 步失败,库存被扣了,但订单没支付,相当于白送。

这就是经典的分布式事务一致性问题。单体时代靠数据库事务一把梭,现在?呵呵。


尝试方案一:两阶段提交(2PC)?想多了

一开始我天真地以为可以用 XA 协议搞个 2PC。查了资料,Python 生态里支持 XA 的 ORM 几乎没有,Django 和 SQLAlchemy 都不原生支持。就算硬上,性能也感人——每个事务都要两次网络往返,延迟直接翻倍。

而且我们用的是 MySQL + Redis 混合存储,XA 根本跨不了 Redis。运维大哥听说我想上 2PC,直接回了个表情包:“你清醒一点”。


尝试方案二:本地消息表 + 轮询?土但有效?

后来我在 Stack Overflow 上看到有人用“本地消息表”模式:在订单服务里建一张 outbox_message 表,下单时同时插入订单记录和一条待发送的消息(比如 “扣减库存”),然后由后台任务轮询这张表,把消息发给库存服务。

听起来可行,于是我吭哧吭哧写了代码:

# orders/models.py
class Order(models.Model):
    user_id = models.IntegerField()
    status = models.CharField(max_length=20)

class OutboxMessage(models.Model):
    event_type = models.CharField(max_length=50)  # e.g., "inventory_deduct"
    payload = models.JSONField()
    processed = models.BooleanField(default=False)
    created_at = models.DateTimeField(auto_now_add=True)

下单逻辑变成:

with transaction.atomic():
    order = Order.objects.create(user_id=123, status="pending")
    OutboxMessage.objects.create(
        event_type="inventory_deduct",
        payload={"order_id": order.id, "sku_id": 456, "qty": 1}
    )

然后有个 Celery 任务每 5 秒扫一次未处理的消息,发 HTTP 请求到库存服务。

结果? 第一天测试就翻车了。

问题出在:如果库存服务返回 500,消息会被重试,但万一重试期间库存已经扣了(只是响应丢了),就会重复扣减!幂等性没做好,直接导致负库存。测试妹子跑来质问:“你们这系统是打算让用户白拿东西吗?”

我当场石化,赶紧加了个去重表:

class DedupKey(models.Model):
    key = models.CharField(max_length=100, unique=True)  # e.g., f"deduct_{order_id}"

每次处理消息前先查这个表,存在就跳过。但这样一来,表越来越大,查询变慢,而且还要定期清理——运维又开始皱眉了。


转折点:拥抱消息队列 + 最终一致性

周五晚上那条企业微信消息之后,我痛定思痛,决定放弃“强一致”的幻想,转向最终一致性。毕竟电商场景下,允许短暂不一致(比如几秒内库存没扣),但不能容忍数据错乱。

我们公司 Kafka 已经铺好了,于是决定用 可靠消息最终一致性 模式。

核心思路:

  1. 订单服务在本地事务中创建订单 + 发送 Kafka 消息(使用事务性生产者)
  2. Kafka 保证消息至少投递一次
  3. 库存服务消费消息,执行扣库存,并保证幂等

关键点来了:如何保证“本地事务 + 发消息”原子性?

答案是:Kafka 的事务性生产者(Transactional Producer)。

Python 里可以用 confluent-kafka 库实现:

from confluent_kafka import Producer

producer = Producer({
    'bootstrap.servers': 'kafka:9092',
    'transactional.id': 'order-service-tx-' + str(uuid.uuid4()),
})

def create_order_and_send_message(user_id, sku_id, qty):
    producer.init_transactions()
    try:
        producer.begin_transaction()
        
        # 1. 创建订单(假设用 Django ORM)
        order = Order.objects.create(user_id=user_id, status="pending")
        
        # 2. 发送消息
        message = {
            "order_id": order.id,
            "sku_id": sku_id,
            "qty": qty
        }
        producer.produce(
            topic='inventory-deduct',
            key=str(order.id),
            value=json.dumps(message).encode('utf-8')
        )
        
        # 3. 提交 Kafka 事务(同时提交本地 DB 事务)
        producer.commit_transaction()
        return order
    except Exception as e:
        producer.abort_transaction()
        raise e

⚠️ 注意:这里假设你的数据库事务和 Kafka 事务能协调——实际上我们用了 两阶段提交的变种:先写 DB,再发 Kafka,但如果 Kafka 失败,就回滚 DB。这依赖于 Kafka 事务的原子性。

但 Python 的 GIL 和异步特性让这事有点 tricky。我们最终选择在业务逻辑层手动控制:先写 DB,再发消息,如果发消息失败,就补偿回滚订单状态(比如标记为“创建失败”)。虽然不是 100% 原子,但在重试机制下,最终能一致。


幂等性:防重复扣库存的终极防线

库存服务那边,我死死记住了“幂等”两个字。每个扣库存请求都带一个唯一 ID(比如 order_id),库存服务用 Redis 做去重:

# inventory/views.py
def deduct_inventory(request):
    order_id = request.data['order_id']
    sku_id = request.data['sku_id']
    qty = request.data['qty']
    
    dedup_key = f"deduct_{order_id}"
    
    # 用 SETNX 实现分布式锁 + 去重
    if redis_client.set(dedup_key, "1", nx=True, ex=3600):
        # 执行扣库存
        stock = Stock.objects.select_for_update().get(sku_id=sku_id)
        if stock.quantity >= qty:
            stock.quantity -= qty
            stock.save()
            return JsonResponse({"status": "ok"})
        else:
            return JsonResponse({"error": "insufficient_stock"}, status=400)
    else:
        # 已处理过,直接返回成功(幂等)
        return JsonResponse({"status": "ok"})

这样即使消息重复投递,也不会多次扣库存。测试妹子终于露出了满意的笑容。


工具链:让分布式事务不再裸奔

光有代码不够,还得有配套工具:

工具 用途 我们的选型
Kafka 可靠消息传递 Confluent Kafka
Sentry 异常监控 自建 Sentry
Prometheus + Grafana 消息积压监控 监控 lag 指标
Flower Celery 任务面板 用于兜底补偿任务

特别提一句:我们还加了个补偿任务,每天凌晨扫描“长时间未完成的订单”,如果发现订单创建成功但库存未扣,就手动触发补偿流程。虽然丑,但救命。


教训与心得

  1. 不要追求强一致:除非是银行转账,否则最终一致性足够用。省下的性能和复杂度,够你多喝十杯咖啡。
  2. 幂等性是生命线:所有下游服务必须幂等,这是分布式系统的铁律。
  3. 可观测性比代码更重要:没监控的分布式系统就像盲人开车。我们上线后第一时间接入了 Kafka lag 监控,避免消息堆积。
  4. Python 不是银弹:虽然我们用 Python 写后端,但分布式事务的难点不在语言,而在设计。Go 或 Java 也一样会踩坑。
  5. 试用期别怕背锅:这次搞定了,领导说“小伙子可以”,转正应该稳了 😅

最后

写这篇文章的时候,已经是凌晨 2 点。窗外安静得只剩键盘声。虽然过程曲折,但看着 Grafana 上平稳的消息消费曲线,心里还是有点小得意。

分布式事务没有银弹,但只要理解业务容忍度、善用工具、做好幂等和监控,就能在“不完美”中找到平衡。

对了,如果你也在试用期,正在被分布式事务折磨——别慌,你不是一个人。欢迎评论区一起吐槽,或者分享你的踩坑故事。

(P.S. 产品经理刚刚又在群里@我,说要加个“预售定金膨胀”功能……我先去续杯咖啡了。)

评论 0

最热最新
暂无评论
MQ堵车了Lv.1
0
影响力
0
文章
0
粉丝