分布式事务解决方案:那些年我踩过的坑,以及怎么爬出来的

淡然如水
2025-12-17 14:23
阅读 1628

大家好,我是老张,一个在某“稳定”国企混了三年半的后端程序员。每天朝九晚五,双休雷打不动,加班?不存在的——我们领导信奉“可持续摸鱼才是高效开发”。不过最近有点心慌,眼瞅着同届进来的同事要么升职加薪,要么跳槽涨薪50%,而我还在用Python写内部审批系统,心里默默盘算:是不是该换个环境了?

为了给自己简历加点硬核内容,最近狠心啃下了分布式事务这块硬骨头。说起来也是被逼的——去年我们部门接了个“数字化转型重点工程”,要把原来单体架构的报销系统拆成微服务。结果刚上线测试环境,财务那边就炸了:“怎么同一个报销单,钱扣了两次?”、“流程走到一半,数据没了?”……那一刻,我真的想把电脑屏幕砸了。

今天这篇不是教科书式的理论科普,而是实打实的踩坑经验分享。我会结合我们项目中 Python + Spring Boot 混合架构的真实场景,聊聊分布式事务到底该怎么搞,以及为什么有些方案听起来很美,用起来却能让你凌晨三点还在群里@运维。


一、问题来了:微服务拆分后的“一致性地狱”

先说背景。原来的报销系统是典型的单体应用,Spring Boot + MySQL,所有操作在一个事务里搞定。但领导非要“上云原生”,于是我们把系统拆成了三个服务:

  • user-service(用户信息,Python + Flask)
  • finance-service(资金处理,Java + Spring Boot)
  • workflow-service(审批流,Spring Boot)

拆完之后,一个完整的报销流程涉及跨三个服务的数据库操作。比如:

  1. 用户提交报销单(写 user-service 的 DB)
  2. 财务冻结预算(写 finance-service 的 DB)
  3. 启动审批流(写 workflow-service 的 DB)

理想很丰满,现实很骨感。第一次联调时,只要中间任何一个服务挂了,整个流程就处于“薛定谔状态”——有的数据写了,有的没写,用户以为提交成功了,财务却查不到记录。测试妹子直接在群里发了个“裂开”的表情包。

这时候我才意识到:本地事务已死,分布式事务当立


二、方案选型:别被PPT忽悠了

网上一搜,分布式事务方案一大堆:2PC、3PC、TCC、Saga、MQ事务消息、Seata、Atomikos……看起来都很牛,但真要落地,得考虑团队技术栈、运维成本、性能要求。

我们团队的情况:

  • Java 系(Spring Boot)主力,但 Python 服务也有重要模块
  • 运维只支持 K8s,不接受额外中间件(除非你写个 YAML 让他 apply)
  • 领导要求“不能有脏数据”,但允许最终一致性(毕竟不是银行)

综合评估后,我们排除了几个选项:

  • 2PC/3PC:强一致性,但性能差、锁粒度大,而且 Atomikos 对 Python 支持几乎为零,pass。
  • TCC:需要每个业务写 Try/Confirm/Cancel 三套逻辑,对我们这种小团队来说维护成本太高,尤其 Python 服务很难统一规范。
  • 纯 MQ 事务消息:依赖 RocketMQ 或 Kafka,但我们的 K8s 集群只部署了 RabbitMQ,改造成本高。

最后,我们决定采用 Saga 模式 + 补偿机制 + 本地消息表 的混合方案。理由很简单:

  • Saga 天然支持最终一致性,适合我们的业务场景
  • 本地消息表实现简单,不依赖外部中间件
  • Python 和 Java 都能轻松实现,代码风格差异不影响核心逻辑

当然,这个决策后来也让我吃了不少苦头……


3. 实战:本地消息表 + Saga 的落地细节

3.1 核心思路

每个服务在执行本地事务的同时,往自己的 outbox_message 表里插入一条待发送的消息。然后由一个独立的异步任务(我们叫它 Message Dispatcher)轮询这张表,把消息发给下游服务。

如果下游服务处理失败,它会返回错误码,上游服务收到后触发补偿操作(比如 undo 冻结预算)。

关键点:本地事务和消息插入必须在一个 DB 事务里完成,确保“要么都成功,要么都失败”。

3.2 Python 侧实现(user-service)

我们在 Flask 里加了一张表:

CREATE TABLE outbox_message (
    id BIGINT AUTO_INCREMENT PRIMARY KEY,
    message_id VARCHAR(64) NOT NULL UNIQUE,
    target_service VARCHAR(32) NOT NULL, -- 'finance', 'workflow'
    payload JSON NOT NULL,
    status ENUM('pending', 'sent', 'failed') DEFAULT 'pending',
    created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
    retry_count INT DEFAULT 0
);

业务代码里这样用:

# user_service/api/expense.py
from .models import db, OutboxMessage
import uuid

def submit_expense(user_id, amount, desc):
    try:
        # 1. 本地业务逻辑:创建报销单
        expense = Expense(user_id=user_id, amount=amount, desc=desc)
        db.session.add(expense)
        
        # 2. 插入消息到 outbox(同一个事务!)
        msg = OutboxMessage(
            message_id=str(uuid.uuid4()),
            target_service="finance",
            payload={"expense_id": expense.id, "amount": amount}
        )
        db.session.add(msg)
        
        db.session.commit()  # 关键:一起提交!
        return {"status": "submitted"}
    except Exception as e:
        db.session.rollback()
        logger.error(f"Submit failed: {e}")
        raise

这里有个坑:一开始我们把消息发送放在 try 块里直接发 HTTP 请求,结果网络抖动导致事务回滚,但消息其实已经发出去了(因为 HTTP 不在 DB 事务里)。后来改成只写 DB,异步发,才解决。

3.3 Spring Boot 侧实现(finance-service)

Java 这边类似,但用了 JPA:

// FinanceService.java
@Transactional
public void freezeBudget(String expenseId, BigDecimal amount) {
    // 1. 本地操作:冻结预算
    Budget budget = budgetRepository.findByUserId(...);
    budget.setFrozen(budget.getFrozen().add(amount));
    budgetRepository.save(budget);
    
    // 2. 插入本地消息(准备通知 workflow)
    OutboxMessage msg = OutboxMessage.builder()
        .messageId(UUID.randomUUID().toString())
        .targetService("workflow")
        .payload(Map.of("expenseId", expenseId))
        .build();
    outboxRepository.save(msg);
}

补偿逻辑也很简单:

// 如果 workflow 通知失败,调用此方法
@Transactional
public void cancelFreeze(String expenseId) {
    Budget budget = ...;
    budget.setFrozen(budget.getFrozen().subtract(originalAmount));
    budgetRepository.save(budget);
    
    // 标记原消息为已补偿
    outboxRepository.markCompensated(expenseId);
}

3.4 异步消息分发器(两边都要跑)

我们用 Celery(Python)和 Spring Scheduler(Java)分别实现了轮询任务:

# Python 版 dispatcher
@app.task
def dispatch_messages():
    pending_msgs = OutboxMessage.query.filter_by(status='pending').limit(100).all()
    for msg in pending_msgs:
        try:
            response = requests.post(
                f"http://{msg.target_service}/api/events",
                json=msg.payload,
                timeout=5
            )
            if response.status_code == 200:
                msg.status = 'sent'
            else:
                msg.retry_count += 1
                if msg.retry_count > 3:
                    msg.status = 'failed'
                    alert_ops("Message failed after 3 retries!")
        except Exception as e:
            logger.error(f"Dispatch error: {e}")
        finally:
            db.session.commit()

Java 那边类似,就不贴了。


4. 踩过的那些坑(血泪史)

坑1:消息重复消费

下游服务可能收到重复消息(比如网络超时重试)。我们一开始没做幂等,结果 finance-service 把预算冻结了两次!

解决方案:每个服务维护一个 processed_message_id 表,收到消息先查是否处理过。

CREATE TABLE processed_message (
    message_id VARCHAR(64) PRIMARY KEY,
    processed_at TIMESTAMP
);

坑2:补偿操作自己也失败了

有一次 finance-service 在执行 cancelFreeze 时 DB 连接池爆了,补偿没成功,数据还是不一致。

解决方案:补偿操作也要加 retry + alert,并且人工介入兜底。我们写了个脚本,每天凌晨扫一遍 outbox_message 里 status=failed 的记录,生成工单给运维。

坑3:Python 和 Java 的 JSON 序列化不一致

Python 用的是 json.dumps,默认把 True 转成 true;Java 用 Jackson,但某个字段用了自定义序列化器,结果变成了 "TRUE"。下游解析直接报错。

解决方案:统一接口契约,用 OpenAPI 3.0 定义 payload 结构,两边生成 client code。虽然麻烦,但值了。

坑4:本地消息表膨胀

三个月后,outbox_message 表涨到了 2000 万行,查询慢得像蜗牛。

解决方案

  • 加索引:(status, created_at)
  • 每天归档成功/失败的消息到历史表
  • 设置 TTL:超过7天的消息自动清理

5. 性能与监控:别让事务拖垮系统

我们压测时发现,加了本地消息表后,QPS 从 1200 降到了 900。主要瓶颈在 DB 写消息。

优化措施:

  • 消息表用单独的 DB 实例(我们用了 TiDB,兼容 MySQL 协议)
  • 批量插入消息(Dispatcher 一次取 100 条,但业务端还是单条写)
  • 异步任务并发数控制在 CPU 核数的 2 倍

监控方面,我们在 Grafana 上建了几个关键看板:

  • 消息积压数量(pending > 100 就告警)
  • 补偿操作触发次数(突然增多说明下游不稳定)
  • 平均事务耗时(超过 500ms 要排查)

6. 为什么没选 Seata?

肯定有人问:你们怎么不用 Seata?AT 模式不是更简单?

原因有三:

  1. 多语言支持弱:Seata 的 Python 客户端基本没人维护,文档少得可怜
  2. 运维复杂:要部署 TC(Transaction Coordinator),在 K8s 里跑 StatefulSet,运维老大直接摇头
  3. 性能损耗:AT 模式要解析 SQL、生成 undo log,对高频交易场景不太友好

对我们这种“求稳不求快”的国企项目,简单可控比炫技更重要


7. 心得体会:分布式事务的本质是妥协

搞完这个项目,我最大的感悟是:没有银弹,只有权衡

  • 强一致性?成本太高,我们业务扛不住。
  • 最终一致性?得接受短暂的不一致,还得有兜底方案。
  • 自研方案?灵活但维护累;开源方案?省事但可能水土不服。

如果你也在做类似改造,我的建议是:

  1. 先画清楚业务边界,哪些操作必须原子,哪些可以异步
  2. 优先考虑本地消息表 + Saga,80% 场景够用
  3. 务必做好幂等、监控、人工干预入口
  4. 别信“完全自动化”,关键时刻还是得人肉救火

最后:关于跳槽的小心思

说实话,折腾完这套分布式事务方案,我简历上终于能写“主导设计高可用分布式系统”了(笑)。上周面试了一家云厂商,聊到 Saga 模式时,面试官眼睛都亮了。

不过话说回来,在国企虽然技术栈旧点,但胜在生活平衡。现在每天下班还能陪老婆追剧,周末带娃去公园——这不比996香?

当然,如果新 offer 给到位,我也不介意再卷一把。毕竟,程序员的终极梦想,不就是用技术换自由吗


附:关键配置参数对比表

方案 一致性级别 性能影响 多语言支持 运维复杂度 适用场景
2PC 强一致 高(锁表) 银行转账
TCC 强一致 电商下单
Saga + 本地消息表 最终一致 优秀 内部系统/非核心交易
Seata AT 强一致 Java为主 Spring Cloud Alibaba 生态

如果你觉得这篇文章对你有帮助,欢迎点赞、转发,或者在评论区吐槽你踩过的分布式事务坑。咱们下期聊聊“K8s 里怎么优雅地做服务熔断”——又是一个让我掉头发的话题 😅

评论 0

最热最新
暂无评论
淡然如水Lv.1
0
影响力
0
文章
0
粉丝