分布式事务解决方案:保障淘宝客APP资金结算的数据强一致性

发布时间:2026/7/25 20:25:51
分布式事务解决方案:保障淘宝客APP资金结算的数据强一致性 分布式事务解决方案保障淘宝客APP资金结算的数据强一致性又见面了我是高佣返利省赚客APP研发者微赚在淘客返利业务中资金结算的准确性是平台的生命线。当用户订单确认收货后系统需要同时完成“订单状态更新”、“佣金计算入账”和“用户余额增加”三个操作。在微服务架构下这三个操作分散在不同的数据库实例中传统的本地Transactional注解失效。一旦某个环节失败极易导致“订单已结算但钱没到账”或“钱扣了但订单未更新”的严重资损事故。为了保障资金数据的强一致性或在极高并发下的最终一致性省赚客APP研发团队深入实践了TCCTry-Confirm-Cancel模式与可靠消息最终一致性方案。TCC模式的三阶段提交实战对于核心资金链路我们放弃了性能较差的XA协议转而采用TCC模式。TCC将事务拆分为Try尝试、Confirm确认、Cancel取消三个阶段。Try阶段进行资源预留和检查Confirm阶段执行真实业务Cancel阶段回滚预留资源。这种模式避免了长事务锁表显著提升了并发性能。packagejuwatech.cn.provinceearn.finance.tcc.service;importjuwatech.cn.provinceearn.finance.entity.UserWallet;importjuwatech.cn.provinceearn.finance.repository.WalletRepository;importjuwatech.cn.provinceearn.common.exception.InsufficientBalanceException;importorg.springframework.stereotype.Service;importorg.springframework.transaction.annotation.Transactional;importlombok.extern.slf4j.Slf4j;importjava.math.BigDecimal;/** * 用户钱包服务的TCC实现 * 负责在分布式事务中管理余额的冻结与扣减 */Slf4jServicepublicclassWalletTccService{privatefinalWalletRepositorywalletRepository;publicWalletTccService(WalletRepositorywalletRepository){this.walletRepositorywalletRepository;}/** * Try阶段检查余额并冻结资金 * 此时资金并未真正扣除只是标记为“冻结”状态防止被其他事务使用 */TransactionalpublicbooleantryDeduct(StringuserId,BigDecimalamount,StringbizId){UserWalletwalletwalletRepository.findByUserId(userId);if(walletnull||wallet.getAvailableBalance().compareTo(amount)0){thrownewInsufficientBalanceException(Insufficient balance for user: userId);}// 记录冻结流水幂等性检查if(walletRepository.existsFreezeRecord(bizId)){returntrue;// 已经冻结过直接返回成功}// 执行冻结可用余额减少冻结余额增加wallet.freezeAmount(amount);walletRepository.save(wallet);// 插入冻结记录表用于Confirm/Cancel时查找walletRepository.saveFreezeRecord(bizId,userId,amount);log.info(Try phase success: frozen {} for user {},amount,userId);returntrue;}/** * Confirm阶段真正扣减资金 * 此阶段必须保证成功若失败需人工介入或无限重试 */TransactionalpublicvoidconfirmDeduct(StringbizId){varrecordwalletRepository.findFreezeRecord(bizId);if(recordnull){thrownewIllegalStateException(Freeze record not found for bizId: bizId);}// 将冻结余额真正扣除清除冻结标记walletRepository.deductFrozenAmount(record.getUserId(),record.getAmount());walletRepository.markFreezeRecordAsConfirmed(bizId);log.info(Confirm phase success: deducted for bizId {},bizId);}/** * Cancel阶段回滚冻结资金 * 将冻结的余额解冻恢复到可用余额 */TransactionalpublicvoidcancelDeduct(StringbizId){varrecordwalletRepository.findFreezeRecord(bizId);if(recordnull){return;// 幂等处理}// 解冻冻结余额减少可用余额增加walletRepository.unfreezeAmount(record.getUserId(),record.getAmount());walletRepository.markFreezeRecordAsCancelled(bizId);log.info(Cancel phase success: unfrozen for bizId {},bizId);}}基于本地消息表的最终一致性方案对于非实时强一致性的场景如积分赠送、通知发送我们采用了“本地消息表定时任务”的方案。该方案的核心思想是将“业务数据写入”与“消息发送”放在同一个本地数据库事务中确保消息一定被持久化。随后由独立的后台线程轮询消息表将消息投递到MQ并在投递成功后更新消息状态。packagejuwatech.cn.provinceearn.finance.message.handler;importjuwatech.cn.provinceearn.finance.entity.DomainEventLog;importjuwatech.cn.provinceearn.finance.repository.EventLogRepository;importjuwatech.cn.provinceearn.common.constant.EventStatus;importorg.apache.rocketmq.spring.core.RocketMQTemplate;importorg.springframework.scheduling.annotation.Scheduled;importorg.springframework.stereotype.Component;importorg.springframework.transaction.annotation.Transactional;importlombok.RequiredArgsConstructor;importjava.time.LocalDateTime;importjava.util.List;/** * 可靠消息投递处理器 * 解决本地事务与MQ发送的原子性问题 */ComponentRequiredArgsConstructorpublicclassReliableMessageSender{privatefinalEventLogRepositoryeventLogRepository;privatefinalRocketMQTemplaterocketMQTemplate;/** * 业务服务调用此方法保存事件必须在业务事务内执行 */TransactionalpublicvoidsaveEvent(StringeventType,StringbusinessId,Stringpayload){DomainEventLogeventLognewDomainEventLog();eventLog.setEventType(eventType);eventLog.setBusinessId(businessId);eventLog.setPayload(payload);eventLog.setStatus(EventStatus.PENDING);eventLog.setCreateTime(LocalDateTime.now());eventLog.setRetryCount(0);// 与业务数据同库同事务保存eventLogRepository.save(eventLog);}/** * 定时任务扫描待发送的消息并投递 * 频率不宜过高避免对数据库造成压力通常设置为每秒或每5秒一次 */Scheduled(fixedRate2000)publicvoidscanAndSend(){// 只查询超过3秒未发送的消息避免干扰刚写入的事务ListDomainEventLogpendingEventseventLogRepository.findPendingEvents(LocalDateTime.now().minusSeconds(3));for(DomainEventLogevent:pendingEvents){try{// 发送消息到MQrocketMQTemplate.syncSend(TOPIC_FINANCE_SETTLEMENT,event.getPayload());// 发送成功更新状态event.setStatus(EventStatus.SENT);event.setUpdateTime(LocalDateTime.now());eventLogRepository.save(event);}catch(Exceptione){// 发送失败增加重试次数event.setRetryCount(event.getRetryCount()1);if(event.getRetryCount()10){event.setStatus(EventStatus.FAILED);log.error(Message failed after 10 retries: {},event.getBusinessId(),e);}else{event.setUpdateTime(LocalDateTime.now());}eventLogRepository.save(event);}}}}防重与补偿机制分布式环境下网络抖动可能导致Confirm或消息重复投递。因此所有下游服务必须具备幂等性。我们在数据库中建立了唯一的业务流水号索引利用数据库的唯一约束来防止重复入账。同时针对TCC模式中Confirm阶段可能出现的长期阻塞我们设计了自动对账补偿任务每日凌晨扫描异常状态的分布式事务自动触发修复或报警人工处理。packagejuwatech.cn.provinceearn.finance.compensation.job;importjuwatech.cn.provinceearn.finance.repository.SettlementOrderRepository;importorg.springframework.batch.core.Job;importorg.springframework.batch.core.Step;importorg.springframework.batch.core.configuration.annotation.JobBuilderFactory;importorg.springframework.batch.core.configuration.annotation.StepBuilderFactory;importorg.springframework.batch.item.database.JdbcPagingItemReader;importorg.springframework.context.annotation.Bean;importorg.springframework.context.annotation.Configuration;importjavax.sql.DataSource;importjava.util.HashMap;importjava.util.Map;/** * 每日对账补偿任务配置 * 扫描状态不一致的结算单进行修复 */ConfigurationpublicclassCompensationJobConfig{privatefinalJobBuilderFactoryjobBuilderFactory;privatefinalStepBuilderFactorystepBuilderFactory;privatefinalDataSourcedataSource;publicCompensationJobConfig(JobBuilderFactoryjobBuilderFactory,StepBuilderFactorystepBuilderFactory,DataSourcedataSource){this.jobBuilderFactoryjobBuilderFactory;this.stepBuilderFactorystepBuilderFactory;this.dataSourcedataSource;}BeanpublicJobreconciliationJob(){returnjobBuilderFactory.get(reconciliationJob).start(reconciliationStep()).build();}BeanpublicStepreconciliationStep(){returnstepBuilderFactory.get(reconciliationStep).MapString,Object,Voidchunk(100).reader(abnormalOrderReader()).processor(item-{// 执行补偿逻辑根据订单状态强制修正资金流// 例如订单已结算但钱包未入账则手动触发入账compensateSettlement(item);returnnull;}).writer(list-{})// 无需写入补偿逻辑在processor中完成.build();}privateJdbcPagingItemReaderMapString,ObjectabnormalOrderReader(){JdbcPagingItemReaderMapString,ObjectreadernewJdbcPagingItemReader();reader.setDataSource(dataSource);reader.setFetchSize(100);reader.setRowMapper((rs,rowNum)-{MapString,ObjectmapnewHashMap();map.put(orderId,rs.getString(order_id));map.put(status,rs.getString(status));returnmap;});// 查询SQL查找状态不匹配的异常数据reader.setSelectStatement(SELECT order_id, status FROM t_settlement_order WHERE status PARTIAL_SUCCESS);returnreader;}privatevoidcompensateSettlement(MapString,Objectitem){// 具体的补偿业务逻辑实现System.out.println(Compensating order: item.get(orderId));}}结语资金安全无小事。通过TCC模式处理核心链路的强一致性需求配合本地消息表实现高吞吐场景下的最终一致性再辅以严密的幂等设计与自动对账补偿省赚客APP构建了一套坚不可摧的资金结算体系。这套架构不仅经受住了双11流量洪峰的考验更确保了每一分返利都准确无误地到达用户手中。本文著作权归 省赚客app 研发团队转载请注明出处