金融交易系统的数据一致性从本地事务到TCC再到Saga的落地复盘一、背景与问题金融交易系统对数据一致性有着天然的高要求——一笔转账涉及扣款与加款两个操作必须要么同时成功要么同时失败不存在扣了款但没加款的中间态。单体架构下本地事务可以轻松保证ACID但微服务拆分后扣款与加款分属不同服务本地事务的边界被打破。本文复盘某支付平台从本地事务到分布式事务的演进路径先尝试TCC方案并遭遇空回滚与悬挂问题最终转向Saga模式并建立完善的补偿与监控体系。二、架构设计概览分布式事务方案的选择取决于业务语义——转账场景要求强一致性不允许出现不一致的中间态而TCC通过Try阶段冻结资源、Confirm/Cancel阶段确认或释放来模拟两阶段提交的强一致语义。但TCC的接口膨胀每个操作需三个接口和空回滚/悬挂等问题增加了工程复杂度。最终我们采用混合策略核心转账链路使用TCC保证强一致外围通知与记录服务使用Saga保证最终一致。TCC的核心思想是资源的「预留-确认-释放」三段式管理而非直接操作。冻结余额而非扣减余额确保Cancel阶段有明确的释放对象。三、核心实现细节3.1 TCC接口设计与空回滚防护TCC的每个业务操作需要拆分为Try、Confirm、Cancel三个接口。以扣款服务为例public class DeductTccService { private final AccountRepository accountRepo; private final TccTransactionLogRepository txLogRepo; /** * Try阶段冻结扣款方余额不实际扣减 */ public TccTryResult tryDeduct(String txId, String accountId, BigDecimal amount) { if (txId null || accountId null || amount null) { throw new TccException(invalid parameters for tryDeduct); } if (amount.compareTo(BigDecimal.ZERO) 0) { throw new TccException(deduct amount must be positive); } // 防悬挂检查是否已存在Cancel记录Cancel先于Try到达 TccTransactionLog existingLog txLogRepo.findByTxIdAndAction(txId, CANCEL); if (existingLog ! null) { log.warn([{}] Hanging Try detected: Cancel already executed, reject Try, txId); return TccTryResult.rejected(hanging Try: Cancel already executed); } Account account accountRepo.findByAccountId(accountId); if (account null) { throw new TccException(account not found: accountId); } BigDecimal availableBalance account.getBalance().subtract(account.getFrozenAmount()); if (availableBalance.compareTo(amount) 0) { return TccTryResult.failed(insufficient available balance); } // 冻结金额而非扣减 account.setFrozenAmount(account.getFrozenAmount().add(amount)); accountRepo.save(account); // 记录Try日志用于后续空回滚检测 txLogRepo.save(TccTransactionLog.builder() .txId(txId) .action(TRY) .accountId(accountId) .amount(amount) .status(SUCCESS) .createdAt(Instant.now()) .build()); return TccTryResult.success(accountId, amount); } /** * Confirm阶段确认扣款将冻结金额转为实际扣减 */ public TccConfirmResult confirmDeduct(String txId, String accountId, BigDecimal amount) { TccTransactionLog tryLog txLogRepo.findByTxIdAndAction(txId, TRY); if (tryLog null) { // Try日志不存在 → 可能是空回滚场景Confirm应直接返回成功 log.warn([{}] No Try log found for Confirm, likely empty rollback scenario, txId); return TccConfirmResult.success(no Try log, skip Confirm); } Account account accountRepo.findByAccountId(accountId); if (account null) { throw new TccException(account not found in Confirm: accountId); } // 冻结金额转为实际扣减 account.setBalance(account.getBalance().subtract(amount)); account.setFrozenAmount(account.getFrozenAmount().subtract(amount)); accountRepo.save(account); txLogRepo.updateStatus(txId, TRY, CONFIRMED); return TccConfirmResult.success(accountId, amount); } /** * Cancel阶段释放冻结金额 * 防空回滚即使Try未执行Cancel也必须正常返回 */ public TccCancelResult cancelDeduct(String txId, String accountId, BigDecimal amount) { TccTransactionLog tryLog txLogRepo.findByTxIdAndAction(txId, TRY); if (tryLog null) { // 空回滚Try未执行但Cancel到达必须记录Cancel日志防止后续Try悬挂 log.warn([{}] Empty rollback: Try not executed, record Cancel log to prevent hanging, txId); txLogRepo.save(TccTransactionLog.builder() .txId(txId) .action(CANCEL) .accountId(accountId) .amount(amount) .status(EMPTY_ROLLBACK) .createdAt(Instant.now()) .build()); return TccCancelResult.success(empty rollback completed); } if (tryLog.getStatus().equals(CONFIRMED)) { // Try已确认Cancel不应执行 log.warn([{}] Cancel after Confirm detected, skip Cancel, txId); return TccCancelResult.success(already confirmed, skip Cancel); } Account account accountRepo.findByAccountId(accountId); if (account null) { throw new TccException(account not found in Cancel: accountId); } // 释放冻结金额 account.setFrozenAmount(account.getFrozenAmount().subtract(amount)); accountRepo.save(account); txLogRepo.updateStatus(txId, TRY, CANCELLED); return TccCancelResult.success(accountId, amount); } }3.2 Saga的正向与补偿流程外围服务采用Saga模式每个正向操作都有对应的补偿操作。补偿操作必须幂等——因为补偿可能因超时重试而被多次触发public class TransferSagaCoordinator { private final SagaDefinitionRepository sagaDefRepo; private final SagaInstanceRepository sagaInstRepo; private final ListSagaStepExecutor stepExecutors; /** * 执行Saga正向流程 * 每步执行前记录状态失败时从当前位置开始逆向补偿 */ public SagaResult execute(String sagaType, TransferContext context) { if (context null || context.getTxId() null) { throw new SagaException(invalid transfer context); } SagaDefinition definition sagaDefRepo.findByType(sagaType); if (definition null) { throw new SagaException(saga definition not found: sagaType); } String sagaInstanceId generateSagaId(); SagaInstance instance SagaInstance.builder() .sagaId(sagaInstanceId) .txId(context.getTxId()) .definitionType(sagaType) .status(RUNNING) .currentStep(0) .createdAt(Instant.now()) .build(); sagaInstRepo.save(instance); ListSagaStep steps definition.getSteps(); for (int i 0; i steps.size(); i) { SagaStep step steps.get(i); instance.setCurrentStep(i); instance.setStepStatus(i, EXECUTING); sagaInstRepo.save(instance); try { SagaStepResult result executeStep(step, context); instance.setStepStatus(i, COMPLETED); instance.setStepOutput(i, result.getOutput()); sagaInstRepo.save(instance); } catch (Exception e) { log.error([{}] Saga step {} failed: {}, sagaInstanceId, step.getName(), e.getMessage()); instance.setStepStatus(i, FAILED); instance.setStatus(COMPENSATING); sagaInstRepo.save(instance); // 逆向补偿已完成步骤 compensateReverse(instance, steps, i - 1, context); return SagaResult.failed(sagaInstanceId, e.getMessage()); } } instance.setStatus(COMPLETED); sagaInstRepo.save(instance); return SagaResult.success(sagaInstanceId); } /** * 逆向补偿从失败位置的上一步开始逐步执行补偿操作 */ private void compensateReverse(SagaInstance instance, ListSagaStep steps, int fromStep, TransferContext context) { for (int i fromStep; i 0; i--) { SagaStep step steps.get(i); String compensateAction step.getCompensateAction(); try { SagaStepResult result executeCompensate(compensateAction, context, instance.getStepOutput(i)); instance.setStepStatus(i, COMPENSATED); sagaInstRepo.save(instance); } catch (Exception e) { log.error([{}] Compensation step {} failed: {}, will retry later, instance.getSagaId(), step.getName(), e.getMessage()); instance.setStepStatus(i, COMPENSATE_FAILED); sagaInstRepo.save(instance); // 补偿失败不中断流程标记待人工处理 break; } } } }四、异常场景覆盖与监控告警4.1 异常场景的系统性覆盖分布式事务的异常场景远比本地事务复杂我们建立了覆盖矩阵确保每种异常都有明确处理策略异常场景TCC处理策略Saga处理策略Try超时记录日志等待Cancel触发空回滚标记步骤失败触发补偿Confirm超时重试Confirm幂等—Cancel超时重试Cancel幂等空回滚安全重试补偿幂等空回滚Cancel先于Try记录空回滚日志阻止后续Try悬挂—悬挂Try在Cancel之后检查Cancel日志拒绝Try—网络分区本地事务日志兜底恢复后重试补偿失败标记待人工4.2 监控告警体系public class TccMonitorService { private final MeterRegistry meterRegistry; private final TccTransactionLogRepository txLogRepo; private final AlertService alertService; /** * 定时扫描异常事务长时间未Confirm/Cancel的Try、悬挂、空回滚 */ Scheduled(fixedDelay 30000) public void scanAbnormalTransactions() { Instant threshold Instant.now().minus(Duration.ofMinutes(5)); // 扫描超时未Confirm的Try ListTccTransactionLog pendingTrys txLogRepo.findByStatusAndCreatedAtBefore(SUCCESS, threshold); for (TccTransactionLog logEntry : pendingTrys) { meterRegistry.counter(tcc.try.pending, service, logEntry.getServiceName()).increment(); alertService.sendAlert(AlertLevel.WARNING, TCC Try pending 5min: txId logEntry.getTxId()); } // 扫描悬挂Try日志存在但对应的Cancel已执行 ListTccTransactionLog hangingTrys txLogRepo.findHangingTrys(); for (TccTransactionLog logEntry : hangingTrys) { meterRegistry.counter(tcc.hanging.try).increment(); alertService.sendAlert(AlertLevel.CRITICAL, TCC hanging Try detected: txId logEntry.getTxId()); } // 扫描空回滚统计 long emptyRollbackCount txLogRepo.countByActionAndStatus(CANCEL, EMPTY_ROLLBACK); meterRegistry.gauge(tcc.empty.rollback.count, emptyRollbackCount); if (emptyRollbackCount 100) { alertService.sendAlert(AlertLevel.WARNING, TCC empty rollback count exceeds threshold: emptyRollbackCount); } } }五、总结金融交易系统的分布式一致性是一场「可靠性工程」而非「完美性追求」。复盘结论如下TCC适合强一致核心链路——资源预留的语义天然适配金融场景但空回滚与悬挂的防护是必做项而非选做项Saga适合外围最终一致链路——补偿操作必须幂等补偿失败必须有兜底的人工处理通道混合策略是务实选择——核心链路TCC、外围链路Saga避免一刀切带来的过度复杂度事务日志是一致性的基石——所有异常场景的判断依据来自事务日志日志的可靠性决定整个方案的可靠性监控告警不能事后补——空回滚、悬挂、超时未确认等异常需要实时监控否则问题会沉默累积直到业务对账才发现下一步演进方向引入事务状态机可视化监控面板使运维人员可以直观追踪每笔分布式事务的状态流转探索基于Seata的AT模式在非核心链路的替代方案降低TCC的接口开发成本。