Administrator
发布于 2022-02-16 / 2456 阅读
42

分布式事务:TCC、Saga、本地消息表怎么选

架构评审会上,两个人对着吵了四十分钟

2 月中旬的一次架构评审,议题是"下单链路要不要上 Seata"。后端 A 主张用 AT 模式,理由是几乎不用改业务代码;后端 B 坚持用消息队列做最终一致,理由是"我们承受不了 Seata 挂掉导致下单全挂"。

我最后拍板:全都不用,主链路走本地消息表,退款单独用 TCC。这篇记的是当时那套判断的依据,以及后面落地时踩到的坑。

先说清楚要解决的问题。下单涉及四个服务:

订单服务   创建订单(MySQL order_db)
库存服务   扣减库存(MySQL stock_db)
优惠券服务 核销券(MySQL coupon_db)
积分服务   冻结积分(MySQL point_db)

四个库,四个独立数据源。我们要的是:要么四步都成,要么都不成。

四种方案的真实成本

方案一致性隔离性业务侵入额外组件我们的人日估算
XA / 2PC强一致好,全局锁几乎为零无(MySQL 原生)5 人日
Seata AT最终(读未提交)全局行锁低,加注解Seata TC 集群12 人日
TCC最终业务预留高,写三套方法事务协调器35 人日
Saga最终无隔离中,写补偿编排器或 MQ25 人日
本地消息表最终无隔离中,建表 + 轮询无(用现有 MQ)10 人日

这个表里的人日数是按我们团队三个人、每个服务平均 6 个接口估的,仅供参考。重点是几个容易被忽略的点:

XA 为什么没人用

MySQL 5.7 之后 XA 是完整的,Spring Boot + Atomikos 也能配。但它的代价是资源锁定时间等于整个事务时长。我们实测过一组数据,同样是下单,四个库各一回合:

本地事务(单库):      P99  28 ms,TPS 1840
XA 两阶段(4 个库):   P99 312 ms,TPS  210

TPS 掉到九分之一。原因是 XA PREPARE 之后,各分支的锁要一直持有到 XA COMMIT,而协调者在网络往返、写 redo 等。四分之三的时间花在协调上。而且 MySQL 的 XA 在 5.7 之前有 prepared 状态事务在崩溃后丢失的 bug(8.0 修了),我们不敢赌。

Seata AT 的真实问题

AT 模式用起来确实简单:

@GlobalTransactional(timeoutMills = 30000, name = "create-order")
public void createOrder(CreateOrderCmd cmd) {
    orderMapper.insert(order);
    stockFeign.deduct(cmd.getSkuId(), cmd.getQty());
    couponFeign.consume(cmd.getCouponId());
}

它靠解析 SQL 生成 undo_log,回滚时用反向 SQL 补偿。两个问题:

一是全局锁。AT 的隔离级别是"读未提交 + 全局行锁",每个分支事务在本地提交前要先到 TC 拿全局锁。TC 是单点,我们压测时 TC 挂掉,整个下单链路 100% 失败。这个风险我们不愿意承担。

二是 undo_log 表。每个业务库都要建,并且会写入大量记录。我们订单库一天新增 60 万单,undo_log 一天就是 240 万行(四个分支),需要额外的清理任务。

我们的选择:主链路本地消息表

下单是典型的"本地事务 + 异步确保"场景:订单落库是主动作,扣库存、核销券、冻积分是可以稍后完成、也可以失败重试的动作。

CREATE TABLE t_local_message (
  id            BIGINT UNSIGNED NOT NULL AUTO_INCREMENT,
  biz_key       VARCHAR(64)  NOT NULL COMMENT '业务唯一键,用于幂等',
  topic         VARCHAR(64)  NOT NULL,
  body          TEXT         NOT NULL,
  status        TINYINT      NOT NULL DEFAULT 0 COMMENT '0待发送 1已发送 2已消费 3失败',
  retry_count   INT          NOT NULL DEFAULT 0,
  next_retry_at DATETIME     NOT NULL,
  create_time   DATETIME     NOT NULL DEFAULT CURRENT_TIMESTAMP,
  update_time   DATETIME     NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
  PRIMARY KEY (id),
  UNIQUE KEY uk_biz_key (biz_key),
  KEY idx_status_retry (status, next_retry_at)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;

发送端:消息和业务数据在同一个本地事务里落库,这是整个方案的关键。

@Transactional(rollbackFor = Exception.class)
public Long createOrder(CreateOrderCmd cmd) {
    Order order = buildOrder(cmd);
    orderMapper.insert(order);

    // 同一个本地事务,写消息表
    localMessageMapper.insert(LocalMessage.builder()
            .bizKey("order-created:" + order.getOrderId())
            .topic("ORDER_CREATED")
            .body(toJson(order))
            .status(0)
            .nextRetryAt(LocalDateTime.now())
            .build());

    return order.getOrderId();
}     // 事务提交 → 订单和消息要么都在,要么都不在

然后由一个定时任务把待发送的消息推到 Kafka,并更新状态:

@Scheduled(fixedDelay = 1000)
public void publish() {
    List<LocalMessage> pending = localMessageMapper.selectPending(100);
    for (LocalMessage msg : pending) {
        try {
            kafkaTemplate.send(msg.getTopic(), msg.getBizKey(), msg.getBody())
                    .get(3, TimeUnit.SECONDS);      // 同步等结果,保证发送成功
            localMessageMapper.markSent(msg.getId());
        } catch (Exception e) {
            localMessageMapper.markFailed(msg.getId());   // retry_count+1,指数退避
            log.warn("消息发送失败 id={} retry={}", msg.getId(), msg.getRetryCount());
        }
    }
}

消费端要做的只有一件事:幂等。因为 Kafka 是 at-least-once,消息可能重复。

@KafkaListener(topics = "ORDER_CREATED", groupId = "stock-service")
public void onOrderCreated(ConsumerRecord<String, String> record) {
    OrderCreatedEvent event = parse(record.value());
    Long orderId = event.getOrderId();

    // 幂等:用唯一索引兜底,而不是先查再插
    try {
        stockRecordMapper.insert(StockRecord.builder()
                .bizKey("deduct:" + orderId)      // 唯一索引
                .skuId(event.getSkuId())
                .qty(-event.getQty())
                .build());
    } catch (DuplicateKeyException e) {
        log.info("重复消费,忽略 orderId={}", orderId);
        return;
    }

    stockMapper.deduct(event.getSkuId(), event.getQty());
}

唯一索引做幂等,不是 select 一下再判断。后者在并发下两个事务同时查到"不存在",然后都插入,就扣了两次库存。

退款为什么用 TCC

退款跟下单不一样:它是金额操作,多退一分钱都是事故。而且退款链路短,只有两个服务(订单状态 + 账务余额),值得为它写三套方法。

public interface AccountTccAction {

    @TwoPhaseBusinessAction(name = "refundAction", commitMethod = "confirm", rollbackMethod = "cancel")
    boolean prepare(BusinessActionContext ctx,
                    @BusinessActionContextParameter(paramName = "userId") Long userId,
                    @BusinessActionContextParameter(paramName = "amount") BigDecimal amount);

    boolean confirm(BusinessActionContext ctx);

    boolean cancel(BusinessActionContext ctx);
}
@Service
public class AccountTccActionImpl implements AccountTccAction {

    @Override
    @Transactional
    public boolean prepare(BusinessActionContext ctx, Long userId, BigDecimal amount) {
        // Try:冻结,不是直接扣
        int n = accountMapper.freeze(userId, amount);   // balance_frozen += amount
        if (n == 0) {
            throw new BizException("余额不足");
        }
        // 记录冻结流水,key 用 xid + branchId,用于幂等和防悬挂
        freezeLogMapper.insert(ctx.getXid(), ctx.getBranchId(), userId, amount);
        return true;
    }

    @Override
    @Transactional
    public boolean confirm(BusinessActionContext ctx) {
        Long userId = ctx.getActionContextLong("userId");
        BigDecimal amount = ctx.getActionContextBigDecimal("amount");
        // Confirm:把冻结的钱真正扣掉
        accountMapper.confirmDeduct(userId, amount);    // balance -= amount, balance_frozen -= amount
        return true;
    }

    @Override
    @Transactional
    public boolean cancel(BusinessActionContext ctx) {
        Long userId = ctx.getActionContextLong("userId");
        BigDecimal amount = ctx.getActionContextBigDecimal("amount");
        // Cancel:释放冻结
        accountMapper.unfreeze(userId, amount);
        return true;
    }
}

Try 阶段的核心是预留资源(冻结),不是直接扣款。这样即使 Confirm 失败,钱还在冻结里,不会凭空消失。

TCC 的三个必答题

这三个问题如果没处理好,TCC 上线就是事故。我在 code review 时专门列了检查项。

1. 幂等

网络抖动时 Confirm / Cancel 会被重复调用。上面代码里 freezeLog 表用 (xid, branch_id) 做唯一索引来兜。

@Override
@Transactional
public boolean confirm(BusinessActionContext ctx) {
    // 幂等判断:这条分支已经 confirm 过就直接返回
    if (freezeLogMapper.isConfirmed(ctx.getXid(), ctx.getBranchId())) {
        return true;
    }
    ...
    freezeLogMapper.markConfirmed(ctx.getXid(), ctx.getBranchId());
}

2. 空回滚

场景:Try 因为网络超时根本没执行成功,但协调者认为超时了,直接发起了 Cancel。Cancel 找不到冻结记录,如果直接返回成功,协调者以为回滚完成了。

@Override
@Transactional
public boolean cancel(BusinessActionContext ctx) {
    FreezeLog log = freezeLogMapper.select(ctx.getXid(), ctx.getBranchId());
    if (log == null) {
        // 空回滚:Try 没执行过,插入一条"已回滚"标记,防止后续 Try 迟到执行成功
        freezeLogMapper.insertCancelOnly(ctx.getXid(), ctx.getBranchId());
        return true;
    }
    accountMapper.unfreeze(ctx.getUserId(), ctx.getAmount());
    freezeLogMapper.markCancelled(ctx.getXid(), ctx.getBranchId());
    return true;
}

关键是那句 insertCancelOnly,它留下一个"这条分支已经回滚过了"的痕迹。

3. 悬挂

场景:Try 请求因为网络拥堵迟到了,协调者已经超时并完成了整个 Cancel,这时候 Try 才到。如果直接执行,钱被冻结之后再也没人来释放,永久悬挂。

@Override
@Transactional
public boolean prepare(BusinessActionContext ctx, Long userId, BigDecimal amount) {
    // 防悬挂:如果这条分支已经回滚过了,拒绝执行 Try
    if (freezeLogMapper.isCancelled(ctx.getXid(), ctx.getBranchId())) {
        log.warn("Try 迟到,分支已回滚,拒绝执行 xid={}", ctx.getXid());
        return false;
    }
    ...
}

这三个检查项没一个能靠框架自动完成,全得业务侧写。这也是 TCC 成本高的主要原因——不是写三个方法,是写三个方法各自的正确性。

Saga 我们用在了一个正合适的场景

大促期间的"跨店满减下单"要连调营销、库存、锁价、资格四个服务,链路长,而且中间可能等用户确认,不适合长事务。这个场景我们用了 Saga,而且用的是编排式(orchestration)而非协同式(choreography)

// 用 Camunda 7.16 做编排,BPMN 里画流程,补偿节点是显式定义的
@Saga
public class FlashSaleSaga {

    @StartSaga
    @SagaEventHandler(associationProperty = "orderId")
    public void handle(OrderCreatedEvent event) {
        commandGateway.send(new LockPriceCommand(event.getOrderId()));
    }

    @SagaEventHandler(associationProperty = "orderId")
    public void handle(PriceLockedEvent event) {
        commandGateway.send(new DeductStockCommand(...));
    }

    @SagaEventHandler(associationProperty = "orderId")
    public void handle(StockDeductFailedEvent event) {
        commandGateway.send(new UnlockPriceCommand(event.getOrderId()));   // 补偿
        commandGateway.send(new ReleaseCouponCommand(event.getOrderId()));
    }
}

Saga 最大的问题是没有隔离性:中间状态对外可见。用户可能看到"库存已扣但券没核销"的订单。我们的处理是在订单上加一个 PROCESSING 状态,前端对这个状态做特殊展示("正在确认中"),并且限制这个状态的订单 5 分钟内不能再次操作。

落地之后的数据

下单链路(本地消息表)
  改造人日        9 人日
  下单 P99        31 ms(改造前 28 ms,增加了 3 ms 的消息表插入)
  消息发送延迟    P99 1.2 秒(1 秒轮询 + 发送耗时)
  不一致订单      上线 3 个月,人工干预 0 起
  消息积压告警    2 次(都是下游服务全挂,恢复后自动追平)

退款链路(TCC)
  改造人日        31 人日(比预估少,因为 Account 服务本来就有冻结字段)
  退款 P99        246 ms
  资金差错        0
  空回滚/悬挂     监控显示触发 11 次,均被正确处理

写在后面

现在回头看,《分布式事务:TCC、Saga、本地消息表怎么选》本身不算多难,难的是线上真出问题那十分钟里的判断。经验都是这么来的。

参考