{T}

分布式事务与最终一致性

单体应用里,一个本地事务往往就能把"扣库存、改订单、记流水"一次性做完;但系统拆成多个服务后,订单、库存、支付、积分分别在不同进程甚至不同数据库里,本地事务边界立刻失效。

这时真正要回答的问题不再是"要不要事务",而是:

  • 哪些业务必须强一致
  • 哪些业务允许短暂不一致
  • 出现失败时靠什么补偿和收敛

所以分布式事务的核心从来不是追求"所有服务像单机一样同时提交",而是按业务代价选择合适的一致性方案。

分布式事务的核心问题

为什么单机事务在分布式里不够用

一次跨服务操作通常会经历:

  1. 服务 A 先更新自己的数据库
  2. 再远程调用服务 B
  3. 服务 B 也更新自己的数据库
  4. 任何一步超时、失败、重试,都可能让整体状态不一致

典型问题包括:

  • A 成功,B 失败:A 已提交,B 未执行,数据不一致
  • A 调用 B 超时,但 B 实际已经执行成功:A 以为失败触发补偿,B 已完成操作
  • 网络抖动触发重试,B 被执行两次:重复扣款、重复发货等问题

这也是为什么分布式事务一定和幂等、超时、重试、补偿一起讨论。

CAP 定理与 BASE 理论

在分布式系统中,CAP 定理指出无法同时满足以下三个特性:

  • C (Consistency) 一致性:所有节点在同一时间看到的数据是一致的
  • A (Availability) 可用性:每个请求都能在合理时间内得到响应
  • P (Partition tolerance) 分区容错性:网络分区发生时系统仍能继续运行

在分布式环境下,网络分区不可避免,因此必须在 C 和 A 之间做权衡。

BASE 理论是对 CAP 的补充:

  • BA (Basically Available) 基本可用:系统出现故障时,允许损失部分可用性
  • S (Soft State) 软状态:允许系统存在中间状态,该状态不影响系统整体可用性
  • E (Eventually Consistent) 最终一致性:经过一段时间后,所有副本最终达到一致状态

分布式事务解决方案对比

方案一致性强度性能影响实现复杂度适用场景
2PC/XA强一致高(持锁时间长)传统数据库、金融核心系统
TCC最终一致中(需要两次调用)资金转账、库存预占、高一致性业务
本地消息表最终一致电商订单、积分发放、通知推送
事务消息最终一致订单创建、库存扣减、异步通知
Saga最终一致长流程业务、编排型事务

常见方案一:两阶段提交 (2PC/XA)

工作原理

理论上可以用 2PC 做强一致协调:

第一阶段(Prepare):

  1. 协调者向所有参与者发送 Prepare 请求
  2. 参与者执行本地事务,但不提交,只写 undo/redo 日志
  3. 参与者返回准备就绪(Yes)或中止(No)

第二阶段(Commit/Rollback):

  1. 如果所有参与者都返回 Yes,协调者发送 Commit 命令
  2. 如果任一参与者返回 No,协调者发送 Rollback 命令
  3. 参与者执行提交或回滚,并释放资源
java
// XA 事务示例(MySQL + Atomikos)
import com.atomikos.jdbc.AtomikosDataSourceBean;
import javax.transaction.UserTransaction;

public class XATransactionExample {
    
    private AtomikosDataSourceBean orderDataSource;
    private AtomikosDataSourceBean inventoryDataSource;
    
    public void createOrderWithXA(Order order) {
        UserTransaction utx = getUserTransaction();
        try {
            utx.begin();
            
            // 操作订单库
            Connection orderConn = orderDataSource.getConnection();
            PreparedStatement orderStmt = orderConn.prepareStatement(
                "INSERT INTO orders (order_no, amount) VALUES (?, ?)"
            );
            orderStmt.setString(1, order.getOrderNo());
            orderStmt.setBigDecimal(2, order.getAmount());
            orderStmt.executeUpdate();
            
            // 操作库存库
            Connection invConn = inventoryDataSource.getConnection();
            PreparedStatement invStmt = invConn.prepareStatement(
                "UPDATE inventory SET stock = stock - ? WHERE product_id = ?"
            );
            invStmt.setInt(1, order.getQuantity());
            invStmt.setLong(2, order.getProductId());
            invStmt.executeUpdate();
            
            utx.commit();
        } catch (Exception e) {
            utx.rollback();
            throw new RuntimeException("XA transaction failed", e);
        }
    }
}

2PC 的问题

它的问题也很明显:

  • 协调成本高:需要协调者维护全局事务状态
  • 持锁时间长:Prepare 阶段会锁住资源直到 Commit/Rollback
  • 对可用性不友好:任一参与者故障会阻塞整个事务
  • 单点故障:协调者故障会导致所有参与者处于阻塞状态
  • 数据不一致风险:Commit 阶段部分参与者失败会导致不一致

所以 2PC 更适合少量强一致场景,不适合作为互联网业务默认主线。

常见方案二:本地消息表 + 异步补偿

工作原理

这是业务里最常见、也最实用的一种做法。

基本思路:

  1. 在服务 A 的本地事务里,同时写业务数据和消息表
  2. 本地事务提交成功后,由后台任务把消息投递到 MQ
  3. 服务 B 消费消息并执行业务逻辑
  4. 如果失败,则重试、告警或进入补偿流程

实现示例

步骤 1:创建本地消息表

sql
CREATE TABLE outbox_message (
    id BIGINT PRIMARY KEY AUTO_INCREMENT,
    topic VARCHAR(100) NOT NULL COMMENT '消息主题',
    biz_key VARCHAR(200) NOT NULL COMMENT '业务唯一键',
    payload TEXT NOT NULL COMMENT '消息内容',
    status VARCHAR(20) NOT NULL COMMENT 'NEW-新建, SENT-已发送, FAILED-失败',
    retry_count INT DEFAULT 0 COMMENT '重试次数',
    create_time DATETIME DEFAULT CURRENT_TIMESTAMP,
    update_time DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
    INDEX idx_status_create (status, create_time),
    UNIQUE KEY uk_biz_key (biz_key)
) COMMENT '本地消息表';

步骤 2:订单服务 - 本地事务写入

java
@Service
public class OrderService {
    
    @Autowired
    private OrderRepository orderRepository;
    
    @Autowired
    private OutboxMessageRepository outboxRepository;
    
    @Transactional(rollbackFor = Exception.class)
    public Order createOrder(CreateOrderCommand command) {
        // 1. 创建订单
        Order order = Order.create(command);
        orderRepository.save(order);
        
        // 2. 写入本地消息表(同一个事务)
        OutboxMessage message = OutboxMessage.builder()
            .topic("order-created")
            .bizKey(order.getOrderNo())
            .payload(toJson(order))
            .status("NEW")
            .build();
        outboxRepository.save(message);
        
        return order;
    }
    
    private String toJson(Object obj) {
        try {
            return new ObjectMapper().writeValueAsString(obj);
        } catch (JsonProcessingException e) {
            throw new RuntimeException("JSON serialization failed", e);
        }
    }
}

步骤 3:消息投递任务

java
@Component
public class MessagePublisher {
    
    @Autowired
    private OutboxMessageRepository outboxRepository;
    
    @Autowired
    private RocketMQTemplate rocketMQTemplate;
    
    @Scheduled(fixedDelay = 5000) // 每 5 秒执行一次
    public void publishPendingMessages() {
        // 查询待发送的消息(限制数量避免一次处理太多)
        List<OutboxMessage> messages = outboxRepository
            .findTop100ByStatusOrderByCreateTimeAsc("NEW");
        
        for (OutboxMessage msg : messages) {
            try {
                // 发送到 MQ
                SendResult result = rocketMQTemplate.syncSend(
                    msg.getTopic(), 
                    msg.getPayload()
                );
                
                // 发送成功,更新状态
                if (result.getSendStatus() == SendStatus.SEND_OK) {
                    msg.setStatus("SENT");
                    msg.setUpdateTime(LocalDateTime.now());
                    outboxRepository.save(msg);
                }
            } catch (Exception e) {
                // 发送失败,增加重试次数
                msg.setRetryCount(msg.getRetryCount() + 1);
                if (msg.getRetryCount() >= 5) {
                    msg.setStatus("FAILED");
                    // 发送告警
                    alertService.sendAlert("消息发送失败: " + msg.getBizKey());
                }
                outboxRepository.save(msg);
            }
        }
    }
}

步骤 4:库存服务 - 消费消息

java
@Service
@RocketMQMessageListener(
    topic = "order-created",
    consumerGroup = "inventory-consumer"
)
public class InventoryConsumer implements RocketMQListener<String> {
    
    @Autowired
    private InventoryService inventoryService;
    
    @Autowired
    private DeductionRecordRepository deductionRepository;
    
    @Override
    public void onMessage(String message) {
        Order order = parseOrder(message);
        
        // 幂等检查:是否已处理过
        if (deductionRepository.existsByOrderNo(order.getOrderNo())) {
            log.info("订单已处理,跳过: {}", order.getOrderNo());
            return;
        }
        
        // 扣减库存
        inventoryService.deductStock(
            order.getProductId(), 
            order.getQuantity(),
            order.getOrderNo()
        );
    }
}

优点与代价

它的优点是:

  • 不依赖全局锁:每个服务只需要管理自己的本地事务
  • 可用性更高:任一服务故障不会阻塞其他服务
  • 更符合最终一致性场景:适合大多数互联网业务

它的代价是:

  • 一致性不是瞬时完成:存在时间窗口的不一致
  • 需要消息投递状态管理:需要维护消息表和投递任务
  • 需要补偿和人工兜底机制:失败场景需要人工介入

常见方案三:事务消息 (RocketMQ)

工作原理

RocketMQ 提供了事务消息机制,解决"本地事务与发消息之间的原子性问题":

  1. 发送半消息:消息发送到 Broker,但对消费者不可见
  2. 执行本地事务:执行本地业务逻辑
  3. 提交或回滚
    • 本地事务成功 → 发送 Commit,消息对消费者可见
    • 本地事务失败 → 发送 Rollback,消息被丢弃
  4. 回查机制:如果 Broker 未收到 Commit/Rollback,会回查事务状态
java
@Service
public class OrderTransactionProducer {
    
    @Autowired
    private RocketMQTemplate rocketMQTemplate;
    
    public void createOrderWithTransaction(Order order) {
        // 构建消息
        Message<String> message = MessageBuilder
            .withPayload(toJson(order))
            .build();
        
        // 发送事务消息
        rocketMQTemplate.sendMessageInTransaction(
            "order-group",
            "order-created",
            message,
            order // 传递给本地事务执行的参数
        );
    }
}

@RocketMQTransactionListener(rocketMQTemplateBeanName = "rocketMQTemplate")
class OrderTransactionListener implements RocketMQLocalTransactionListener {
    
    @Autowired
    private OrderService orderService;
    
    @Autowired
    private OrderRepository orderRepository;
    
    @Override
    @Transactional(rollbackFor = Exception.class)
    public RocketMQLocalTransactionState executeLocalTransaction(Message msg, Object arg) {
        try {
            Order order = (Order) arg;
            // 执行本地事务:创建订单
            orderService.createOrder(order);
            
            // 返回提交状态,消息对消费者可见
            return RocketMQLocalTransactionState.COMMIT;
        } catch (Exception e) {
            // 返回回滚状态,消息被丢弃
            return RocketMQLocalTransactionState.ROLLBACK;
        }
    }
    
    @Override
    public RocketMQLocalTransactionState checkLocalTransaction(Message msg) {
        // 回查本地事务状态
        String orderNo = extractOrderNo(msg);
        
        if (orderRepository.existsByOrderNo(orderNo)) {
            // 订单存在,提交消息
            return RocketMQLocalTransactionState.COMMIT;
        } else {
            // 订单不存在,回滚消息
            return RocketMQLocalTransactionState.ROLLBACK;
        }
    }
}

事务消息解决的问题

事务消息解决的是"本地事务与发消息之间的原子性问题",但它并不会自动替你解决:

  • 消费者幂等:消费者必须自己实现幂等处理
  • 消息重试:消费失败后的重试策略
  • 业务补偿:长时间失败后的补偿机制

常见方案四:TCC (Try-Confirm-Cancel)

工作原理

TCC = TryConfirmCancel

它适合对一致性要求更高、且业务可明确拆成预留和确认阶段的场景,例如:

  • 账户余额预冻结
  • 库存预占
  • 优惠券预锁定

三个阶段可以这样理解:

  • Try:先预留资源,不真正完成业务
  • Confirm:所有环节都准备好后正式提交
  • Cancel:任一环节失败时释放预留资源

TCC 实现示例

步骤 1:定义 TCC 接口

java
@LocalTCC
public interface InventoryTccService {
    
    /**
     * Try: 预占库存
     */
    @TwoPhaseBusinessAction(
        name = "prepareDeduct",
        commitMethod = "commit",
        rollbackMethod = "rollback"
    )
    boolean prepareDeduct(
        @BusinessActionContextParameter(paramName = "productId") Long productId,
        @BusinessActionContextParameter(paramName = "quantity") Integer quantity,
        @BusinessActionContextParameter(paramName = "orderNo") String orderNo
    );
    
    /**
     * Confirm: 确认扣减
     */
    boolean commit(BusinessActionContext context);
    
    /**
     * Cancel: 释放预留库存
     */
    boolean rollback(BusinessActionContext context);
}

步骤 2:实现 TCC 服务

java
@Service
public class InventoryTccServiceImpl implements InventoryTccService {
    
    @Autowired
    private InventoryRepository inventoryRepository;
    
    @Autowired
    private InventoryFreezeRepository freezeRepository;
    
    @Override
    @Transactional(rollbackFor = Exception.class)
    public boolean prepareDeduct(Long productId, Integer quantity, String orderNo) {
        // 1. 检查库存是否充足
        Inventory inventory = inventoryRepository.findByProductId(productId);
        if (inventory.getAvailableStock() < quantity) {
            throw new RuntimeException("库存不足");
        }
        
        // 2. 冻结库存(预占)
        inventory.setAvailableStock(inventory.getAvailableStock() - quantity);
        inventory.setFrozenStock(inventory.getFrozenStock() + quantity);
        inventoryRepository.save(inventory);
        
        // 3. 记录冻结明细(用于回滚)
        InventoryFreeze freeze = InventoryFreeze.builder()
            .orderNo(orderNo)
            .productId(productId)
            .quantity(quantity)
            .status("FROZEN")
            .createTime(LocalDateTime.now())
            .build();
        freezeRepository.save(freeze);
        
        return true;
    }
    
    @Override
    @Transactional(rollbackFor = Exception.class)
    public boolean commit(BusinessActionContext context) {
        String orderNo = context.getActionContext("orderNo", String.class);
        Integer quantity = context.getActionContext("quantity", Integer.class);
        Long productId = context.getActionContext("productId", Long.class);
        
        // 1. 删除冻结记录(幂等处理)
        InventoryFreeze freeze = freezeRepository.findByOrderNo(orderNo);
        if (freeze == null) {
            // 已处理过,幂等返回
            return true;
        }
        freezeRepository.delete(freeze);
        
        // 2. 真正扣减冻结库存
        Inventory inventory = inventoryRepository.findByProductId(productId);
        inventory.setFrozenStock(inventory.getFrozenStock() - quantity);
        inventoryRepository.save(inventory);
        
        return true;
    }
    
    @Override
    @Transactional(rollbackFor = Exception.class)
    public boolean rollback(BusinessActionContext context) {
        String orderNo = context.getActionContext("orderNo", String.class);
        Integer quantity = context.getActionContext("quantity", Integer.class);
        Long productId = context.getActionContext("productId", Long.class);
        
        // 1. 查询冻结记录(幂等处理)
        InventoryFreeze freeze = freezeRepository.findByOrderNo(orderNo);
        if (freeze == null) {
            // 已回滚过,幂等返回
            return true;
        }
        
        // 2. 恢复可用库存
        Inventory inventory = inventoryRepository.findByProductId(productId);
        inventory.setAvailableStock(inventory.getAvailableStock() + quantity);
        inventory.setFrozenStock(inventory.getFrozenStock() - quantity);
        inventoryRepository.save(inventory);
        
        // 3. 删除冻结记录
        freezeRepository.delete(freeze);
        
        return true;
    }
}

TCC 的三大问题

TCC 的问题不是理论,而是落地成本高:

1. 幂等问题

问题:Try、Confirm、Cancel 都可能被重复调用

解决:每个阶段都要实现幂等检查

java
@Override
public boolean commit(BusinessActionContext context) {
    String orderNo = context.getActionContext("orderNo", String.class);
    
    // 幂等检查:是否已处理
    InventoryFreeze freeze = freezeRepository.findByOrderNo(orderNo);
    if (freeze == null) {
        log.info("订单已提交,幂等返回: {}", orderNo);
        return true;
    }
    
    // 执行提交逻辑...
}

2. 空回滚问题

问题:Try 未执行成功,但 Cancel 被调用

解决:在 Cancel 阶段检查 Try 是否执行

java
@Override
public boolean rollback(BusinessActionContext context) {
    String orderNo = context.getActionContext("orderNo", String.class);
    
    // 空回滚检查:Try 是否执行过
    InventoryFreeze freeze = freezeRepository.findByOrderNo(orderNo);
    if (freeze == null) {
        log.warn("Try 未执行,空回滚: {}", orderNo);
        // 记录空回滚标记,防止悬挂
        recordEmptyRollback(orderNo);
        return true;
    }
    
    // 执行回滚逻辑...
}

3. 悬挂问题

问题:Cancel 比 Try 先执行,导致 Try 无法执行

解决:在 Try 阶段检查是否已回滚

java
@Override
public boolean prepareDeduct(Long productId, Integer quantity, String orderNo) {
    // 悬挂检查:是否已回滚
    if (hasEmptyRollbackRecord(orderNo)) {
        log.warn("已回滚,拒绝 Try: {}", orderNo);
        return false;
    }
    
    // 执行 Try 逻辑...
}

常见方案五:Saga 模式

工作原理

Saga 模式将长事务拆分为多个本地事务,每个事务都有对应的补偿操作:

  • 正向操作:T1 → T2 → T3 → ... → Tn
  • 补偿操作:T1 补偿 → T2 补偿 → T3 补偿 → ... → Tn 补偿

如果某个步骤失败,按照相反顺序执行补偿操作。

Saga 实现示例(编排模式)

java
@Service
public class OrderSagaService {
    
    @Autowired
    private OrderService orderService;
    
    @Autowired
    private InventoryService inventoryService;
    
    @Autowired
    private PaymentService paymentService;
    
    @Autowired
    private PointService pointService;
    
    public void createOrderSaga(Order order) {
        List<Runnable> compensations = new ArrayList<>();
        
        try {
            // 步骤 1:创建订单
            orderService.createOrder(order);
            compensations.add(() -> orderService.cancelOrder(order.getOrderNo()));
            
            // 步骤 2:扣减库存
            inventoryService.deductStock(order.getProductId(), order.getQuantity());
            compensations.add(() -> inventoryService.restoreStock(
                order.getProductId(), 
                order.getQuantity()
            ));
            
            // 步骤 3:扣款
            paymentService.charge(order.getUserId(), order.getAmount());
            compensations.add(() -> paymentService.refund(
                order.getUserId(), 
                order.getAmount()
            ));
            
            // 步骤 4:发放积分(允许失败,不补偿)
            try {
                pointService.awardPoints(order.getUserId(), order.getPoints());
            } catch (Exception e) {
                log.warn("积分发放失败,不影响主流程", e);
            }
            
        } catch (Exception e) {
            // 执行补偿(按相反顺序)
            Collections.reverse(compensations);
            for (Runnable compensation : compensations) {
                try {
                    compensation.run();
                } catch (Exception ex) {
                    log.error("补偿失败", ex);
                }
            }
            throw new RuntimeException("Saga 执行失败", e);
        }
    }
}

实战场景

场景一:订单、库存、积分的最终一致

用户下单后:

  1. 订单服务先创建订单
  2. 订单服务写本地消息表
  3. 库存服务消费"订单已创建"事件进行扣减
  4. 积分服务消费"订单已支付"事件发放积分

这里更适合最终一致,而不是要求订单、库存、积分在一个全局事务里同时提交。真正的关键在于:

  • 每个消费者都幂等
  • 消息投递失败可重试
  • 超过阈值后有补偿任务和人工告警

状态流转图:

code
订单状态:CREATED → PAID → COMPLETED
库存状态:已扣减
积分状态:待发放 → 已发放

失败场景:
- 库存扣减失败 → 订单取消
- 积分发放失败 → 定时补偿重试

场景二:支付完成后的状态收敛

支付回调最常见的问题是重复通知和乱序通知。

稳妥做法通常是:

  • 订单状态机只允许 INIT -> PAID
  • 支付成功事件按业务单号幂等处理
  • 积分、发票、通知等下游都通过异步事件驱动
  • 如果某个下游失败,不回滚支付结果,而是补偿重试

这个场景说明:核心状态要先收敛,外围动作通过最终一致补齐。

实现要点:

java
@Service
public class PaymentCallbackService {
    
    @Transactional(rollbackFor = Exception.class)
    public void handlePaymentCallback(PaymentCallback callback) {
        // 1. 幂等检查
        if (paymentRecordRepository.existsByPayNo(callback.getPayNo())) {
            log.info("支付回调已处理: {}", callback.getPayNo());
            return;
        }
        
        // 2. 查询订单
        Order order = orderRepository.findByOrderNo(callback.getOrderNo());
        
        // 3. 状态机检查
        if (!order.canTransitionTo(OrderStatus.PAID)) {
            throw new IllegalStateException(
                "订单状态不允许支付: " + order.getStatus()
            );
        }
        
        // 4. 更新订单状态
        order.setStatus(OrderStatus.PAID);
        order.setPayTime(LocalDateTime.now());
        orderRepository.save(order);
        
        // 5. 记录支付流水
        PaymentRecord record = PaymentRecord.builder()
            .payNo(callback.getPayNo())
            .orderNo(order.getOrderNo())
            .amount(callback.getAmount())
            .payTime(callback.getPayTime())
            .build();
        paymentRecordRepository.save(record);
        
        // 6. 发送下游事件(异步)
        eventPublisher.publish("order-paid", order);
    }
}

场景三:优惠券和余额冻结

如果一个链路里同时涉及账户余额冻结和优惠券锁定,且用户感知非常强,可能会考虑用 TCC:

  • Try 阶段先冻结余额、锁券
  • Confirm 阶段真正扣款和核销
  • Cancel 阶段释放冻结资源

但如果业务复杂度没有高到这个程度,贸然上 TCC 往往比问题本身更难维护。

TCC 适用判断标准:

判断维度适合 TCC不适合 TCC
一致性要求强一致,不允许任何时间窗口不一致最终一致,短暂不一致可接受
业务可拆分能明确区分预留和确认阶段业务逻辑简单,无法拆分
失败影响失败会导致严重资金/库存问题失败可补偿,影响可控
开发成本有足够资源实现 TCC 三阶段开发资源有限

排查与治理思路

先判断你缺的是哪一种一致性

不是所有"不一致"都需要全局事务。

先问三个问题:

  1. 这个场景允许多长时间的不一致窗口
  2. 不一致发生后,能不能通过补偿修复
  3. 修复失败后,业务损失是否可接受

如果答案是"可以异步修复",那么优先考虑最终一致;如果答案是"必须立刻一致,否则会产生严重资金或库存风险",再考虑更强的方案。

一致性方案选择决策树

code
是否允许短暂不一致?
├─ 否 → 考虑 2PC/XA(强一致,性能低)
└─ 是
    └─ 业务是否可拆分为 Try-Confirm-Cancel?
        ├─ 是 → 考虑 TCC(高一致性,实现复杂)
        └─ 否
            └─ 是否需要保证消息发送与本地事务的原子性?
                ├─ 是 → 事务消息(推荐 RocketMQ)
                └─ 否 → 本地消息表 + 异步补偿

核心治理动作

分布式事务方案落地时,至少要补齐这些能力:

能力说明实现方式
幂等重复请求、重复消息都不能让结果失控业务单号 + 唯一索引 / 状态机约束
状态机限制非法状态流转,避免重复扣减数据库字段约束 + 业务校验
超时避免远程调用无限等待设置合理的超时时间(如 3-5 秒)
重试只对可恢复错误生效指数退避 + 最大重试次数
补偿失败后有自动修正路径定时任务扫描 + 人工处理入口
告警长时间未收敛的事务必须暴露出来监控系统 + 告警通知

如何避免"假最终一致"

很多系统嘴上说最终一致,实际上只是"先不管"。

真正可用的最终一致至少要有:

能力说明实现方式
明确的收敛目标定义最终一致的状态状态机设计
可观测的中间状态能看到不一致的状态监控面板 + 日志
定时补偿任务自动修正不一致状态定时任务扫描
失败重试上限避免无限重试最大重试次数 + 死信队列
人工处理入口无法自动修复时人工介入管理后台 + 告警通知

没有这些能力,最终一致就会退化成长期脏数据。

分布式事务监控与排查

关键监控指标

指标类型具体指标阈值建议
事务状态进行中的事务数量< 1000
事务状态超时事务数量< 10
事务状态失败事务数量< 5/hour
消息队列消息积压量< 10000
消息队列消息重试次数< 5
补偿任务补偿成功率> 95%
补偿任务人工介入数量< 5/day

常见排查工具

sql
-- 查询长时间未完成的事务
SELECT * FROM outbox_message 
WHERE status = 'NEW' 
  AND create_time < DATE_SUB(NOW(), INTERVAL 10 MINUTE)
ORDER BY create_time;

-- 查询失败的消息
SELECT * FROM outbox_message 
WHERE status = 'FAILED' 
ORDER BY create_time DESC;

-- 查询重试次数过多的消息
SELECT * FROM outbox_message 
WHERE retry_count > 3 
ORDER BY retry_count DESC;

常见误区

  • 一上来就追求全局强一致:结果把系统可用性和复杂度都拖垮
  • 只做消息发送,不做消费幂等和补偿:最后形成重复扣减或长期脏数据
  • 把 TCC 当通用模板,到处套用:忽略业务侵入和维护成本
  • 只靠重试,不做状态机和失败隔离:最终把错误放大
  • 没有中间状态监控,却说系统是最终一致:无法发现和修复不一致

最佳实践总结

选择方案的原则

  1. 能用最终一致就不用强一致:性能和可用性更好
  2. 能用消息表就不用 TCC:实现复杂度更低
  3. 幂等优于锁:幂等是基础能力,锁是补充手段
  4. 监控优于补偿:先发现问题,再解决问题

必须具备的基础能力

  1. 幂等设计:每个接口都要考虑幂等
  2. 状态机设计:限制非法状态流转
  3. 监控告警:实时发现不一致问题
  4. 补偿机制:自动 + 人工两套方案
  5. 文档记录:记录事务流程和异常处理

面试补充

高频面试题

Q1:为什么互联网业务里最终一致性比强一致更常见?

A:因为可用性、吞吐和实现复杂度更可控。强一致(如 2PC)需要长时间持锁,性能差且可用性低;最终一致性允许短暂不一致,通过补偿机制保证最终正确,更适合互联网业务的高并发场景。

Q2:本地消息表解决了什么问题?

A:把业务数据和待发送消息放进同一个本地事务,保证了"业务成功 + 消息发送"的原子性。避免了"业务已提交但消息没发出去"的问题。

Q3:事务消息解决了什么问题?

A:降低"本地事务成功但消息没发出去"的窗口。通过半消息机制,确保本地事务和消息发送的原子性,同时提供了回查机制处理异常情况。

Q4:TCC 适合什么场景?

A:高价值、强一致、可明确拆成预留和确认阶段的业务,如资金转账、库存预占、优惠券锁定等。不适合一般业务,因为实现复杂度很高。

Q5:最终一致性为什么一定要配补偿?

A:因为异步链路天然存在失败和延迟,必须有收敛机制。没有补偿,最终一致就会退化成长期脏数据。

Q6:TCC 的三大问题是什么?如何解决?

A:

  1. 幂等问题:Try、Confirm、Cancel 可能重复调用 → 每个阶段都要实现幂等检查
  2. 空回滚问题:Try 未执行成功但 Cancel 被调用 → Cancel 阶段检查 Try 是否执行
  3. 悬挂问题:Cancel 比 Try 先执行 → Try 阶段检查是否已回滚

Q7:如何选择分布式事务方案?

A:

  1. 需要强一致 → 2PC/XA
  2. 需要高一致且业务可拆分 → TCC
  3. 需要保证消息原子性 → 事务消息
  4. 一般业务 → 本地消息表 + 异步补偿
  5. 长流程业务 → Saga 模式

Q8:分布式事务中如何保证幂等?

A:

  1. 业务唯一单号:每个业务操作有唯一标识
  2. 数据库唯一约束:唯一索引防止重复插入
  3. 状态机约束:只允许合法的状态流转
  4. 去重表:记录已处理的请求 ID
  5. Redis 标记:快速判断是否已处理(需配合持久化)

版本差异(技术原理说明)

维度说明
技术原理分布式一致性/事务/锁/ID 生成等原理与具体版本无关,长期有效
落地选型新项目建议优先使用 Nacos/Redis/Seata 等成熟组件(JDK 17+ 兼容)
Java 版本示例代码基于 JDK 8 编写,JDK 17/21 下语法兼容

本文讲解的分布式系统核心问题与解决方案原理稳定,不随框架版本变化;落地时选用支持 JDK 17/21 的组件版本即可。