SpringBoot中事务内同步处理 + 事务后异步调用外部系统的通用模式示例

📅 2026/7/30 10:59:25
SpringBoot中事务内同步处理 + 事务后异步调用外部系统的通用模式示例
SpringBoot中事务内同步处理 事务后异步调用外部系统的通用模式示例一、解决的核心问题在业务系统中经常遇到这样的需求一个用户操作既要修改本地数据库又要调用外部系统HTTP接口、第三方平台等。直接在一个事务里做这两件事会产生严重问题❌ 错误做法 开启事务 → 写本地数据库 → 调用外部HTTP接口(耗时2-5秒) → 提交事务 问题1长事务 — HTTP超时期间数据库连接和行锁被占用阻塞其他请求 问题2数据不一致 — 外部接口成功了但本地事务提交失败外部状态无法回退 问题3数据丢失 — 外部接口超时/异常导致本地事务回滚之前的写入全部丢失二、解决方案事务提交后异步调用将流程拆分为两阶段阶段1同步事务内校验 写库 注册事务后回调 阶段2异步事务外MQ消费 → 调用外部系统 → 根据结果更新状态 ┌────────────── 阶段1用户请求处理 ──────────────┐ │ │ │ 开始事务 │ │ → 前置校验权限、数据完整性等 │ │ → 准备数据生成编号、组装参数等 │ │ → 写入数据库状态设为处理中 │ │ → 注册 afterCommit 回调发送MQ消息 │ │ 提交事务 │ │ → 触发回调 → MQ消息发出 │ │ │ └────────────────────────────────────────────────────┘ │ (消息队列) │ ▼ ┌────────────── 阶段2异步消费处理 ──────────────┐ │ │ │ MQ Consumer 收到消息 │ │ → 调用外部系统HTTP/RPC │ │ → 成功更新状态为完成 │ │ → 失败更新状态为失败可重试 │ │ → 记录操作日志 │ │ │ └────────────────────────────────────────────────────┘注博客https://blog.csdn.net/badao_liumang_qizhi三、完整通用示例场景假设用户提交一份报告系统需要本地保存报告数据调用外部审批平台创建审批流程审批平台返回审批ID后本地关联保存3.1 实体和枚举// 报告状态枚举publicenumReportStatus{DRAFT(0,草稿),SUBMITTING(1,提交中),// 已提交但审批平台还没响应SUBMITTED(2,已提交),// 审批平台已确认SUBMIT_FAILED(3,提交失败);// 审批平台调用失败privateIntegercode;privateStringdesc;}// 报告实体EntityTable(namereport)publicclassReport{IdprivateIntegerid;privateStringtitle;privateStringcontent;privateIntegerstatus;// 报告状态privateStringapprovalId;// 外部审批平台返回的IDprivateLocalDateTimesubmitTime;}// 操作日志实体记录与外部系统的每次交互EntityTable(namereport_submit_log)publicclassReportSubmitLog{IdprivateIntegerid;privateIntegerreportId;privateStringrequestData;// 发送给外部的数据JSONprivateStringresponseData;// 外部返回的数据JSONprivateStringtransmitFlag;// T成功, F失败privateStringerrorMsg;// 失败原因privateLocalDateTimecreateTime;}3.2 阶段1Service 层处理用户请求ServicepublicclassReportServiceImplimplementsReportService{AutowiredprivateReportRepositoryreportRepository;AutowiredprivateReportSubmitMqSenderreportSubmitMqSender;/** * 提交报告用户直接调用的方法. * 这个方法在事务内执行只做本地数据操作。 */Transactional(rollbackForException.class)OverridepublicResultVoidsubmitReport(IntegerreportId,IntegeruserId){// 第一步前置校验 ReportreportreportRepository.findById(reportId).orElseThrow(()-newBusinessException(报告不存在));if(!ReportStatus.DRAFT.getCode().equals(report.getStatus())){thrownewBusinessException(只有草稿状态的报告才能提交);}// 校验用户权限if(!userId.equals(report.getCreatorId())){thrownewBusinessException(只有创建者才能提交);}// 第二步准备数据、修改本地状态 report.setStatus(ReportStatus.SUBMITTING.getCode());// 设为提交中report.setSubmitTime(LocalDateTime.now());reportRepository.save(report);// 第三步组装外部调用参数 ApprovalCreateParamapprovalParamnewApprovalCreateParam();approvalParam.setTitle(report.getTitle());approvalParam.setContent(report.getContent());approvalParam.setCallbackUrl(http://my-service/api/approval-callback);// 第四步注册事务提交后的回调 // 关键不在事务内发MQ而是注册事务提交成功后的回调TransactionSynchronizationManager.registerSynchronization(newTransactionSynchronization(){OverridepublicvoidafterCommit(){// 这里在事务成功提交后才执行reportSubmitMqSender.send(approvalParam,reportId);}// 如果事务回滚afterCommit 不会被调用MQ消息不会发出});returnResult.success();}}3.3 事务后回调的封装工具类上面直接写匿名类比较繁琐实际项目中通常封装一个收集器/** * 事务后置动作收集器. * 用于收集需要在事务提交后执行的动作通常是发送MQ消息. */publicclassAfterCommitActionCollectorimplementsTransactionSynchronization{privatefinalListRunnableactionsnewArrayList();/** * 添加一个事务提交后要执行的动作. */publicvoidaddAction(Runnableaction){actions.add(action);}OverridepublicvoidafterCommit(){// 事务提交成功后依次执行所有注册的动作for(Runnableaction:actions){try{action.run();}catch(Exceptione){// 记录日志但不抛异常因为本地事务已经提交了log.warn(事务后置动作执行失败,e);}}}}使用方式更简洁// 在事务方法中AfterCommitActionCollectorcollectornewAfterCommitActionCollector();TransactionSynchronizationManager.registerSynchronization(collector);// 可以注册多个动作collector.addAction(()-reportSubmitMqSender.send(approvalParam,reportId));collector.addAction(()-notificationMqSender.send(userId,报告已提交));3.4 MQ 生产者ComponentpublicclassReportSubmitMqSender{AutowiredprivateRabbitTemplaterabbitTemplate;/** * 发送MQ消息. * 消息体包含调用外部系统所需的参数和业务ID. */publicvoidsend(ApprovalCreateParamparam,IntegerreportId){MapString,ObjectmessagenewHashMap();message.put(param,param);message.put(reportId,reportId);log.info(发送报告提交MQ: reportId{},reportId);try{rabbitTemplate.convertAndSend(report.exchange,report.submit,message);}catch(Exceptione){// MQ发送失败的兜底记录日志后续由定时任务扫描提交中超时的记录来补偿log.error(报告提交MQ发送失败: reportId{},reportId,e);}}}3.5 阶段2MQ 消费者ComponentpublicclassReportSubmitMqConsumer{AutowiredprivateReportServicereportService;/** * 消费MQ消息调用外部审批平台. */RabbitListener(queuesreport.submit.queue)publicvoidconsume(MapString,Objectmessage){ApprovalCreateParamparam(ApprovalCreateParam)message.get(param);IntegerreportId(Integer)message.get(reportId);log.info(消费报告提交MQ: reportId{},reportId);// 委托给Service层处理Service中有独立事务reportService.callExternalApproval(param,reportId);}}3.6 阶段2调用外部系统的业务逻辑ServicepublicclassReportServiceImplimplementsReportService{AutowiredprivateReportRepositoryreportRepository;AutowiredprivateReportSubmitLogRepositorysubmitLogRepository;AutowiredprivateApprovalPlatformClientapprovalClient;// HTTP客户端/** * 调用外部审批平台MQ消费后执行. * 使用 REQUIRES_NEW 确保独立事务不受Consumer框架事务影响. */Transactional(propagationPropagation.REQUIRES_NEW,rollbackForException.class)OverridepublicvoidcallExternalApproval(ApprovalCreateParamparam,IntegerreportId){ReportreportreportRepository.findById(reportId).orElse(null);if(reportnull){return;// 数据被删除了直接结束}// 创建操作日志无论成败都要记录ReportSubmitLogsubmitLognewReportSubmitLog();submitLog.setReportId(reportId);submitLog.setRequestData(JsonUtil.toJson(param));submitLog.setCreateTime(LocalDateTime.now());try{// 调用外部系统 ApprovalCreateResultresultapprovalClient.createApproval(param);if(result.isSuccess()){// 成功更新本地状态report.setStatus(ReportStatus.SUBMITTED.getCode());report.setApprovalId(result.getApprovalId());submitLog.setTransmitFlag(T);submitLog.setResponseData(JsonUtil.toJson(result));}else{// 业务失败状态回退report.setStatus(ReportStatus.SUBMIT_FAILED.getCode());submitLog.setTransmitFlag(F);submitLog.setErrorMsg(result.getErrorMsg());}reportRepository.save(report);}catch(Exceptione){// 异常状态回退记录错误report.setStatus(ReportStatus.SUBMIT_FAILED.getCode());reportRepository.save(report);submitLog.setTransmitFlag(F);submitLog.setErrorMsg(e.getMessage());log.warn(调用审批平台失败: reportId{},reportId,e);}finally{// 无论如何都保存操作日志submitLogRepository.save(submitLog);}}}四、为什么用REQUIRES_NEWTransactional(propagationPropagation.REQUIRES_NEW)MQ Consumer 本身可能有事务上下文框架自动管理ACK用REQUIRES_NEW开启独立事务的好处Consumer 框架事务 │ ├── callExternalApproval() 的独立事务 │ → 成功独立提交 │ → 失败独立回滚不影响Consumer的ACK │ └── Consumer 正常结束ACK消息如果不用REQUIRES_NEW业务异常会导致 Consumer 事务回滚 → MQ消息重新入队 → 无限重试。五、失败重试机制定时扫描补偿Scheduled(fixedRate300000)// 每5分钟publicvoidretryFailedSubmissions(){// 查找提交中超过10分钟的记录可能MQ丢失了ListReportstuckReportsreportRepository.findByStatusAndSubmitTimeBefore(ReportStatus.SUBMITTING.getCode(),LocalDateTime.now().minusMinutes(10));for(Reportreport:stuckReports){// 重新发送MQreportSubmitMqSender.send(buildParam(report),report.getId());}}手动重试接口GetMapping(/retry-submit)publicResultVoidretrySubmit(RequestParamIntegerreportId){ReportSubmitLoglastLogsubmitLogRepository.findTopByReportIdOrderByCreateTimeDesc(reportId);if(lastLog!nullF.equals(lastLog.getTransmitFlag())){ApprovalCreateParamparamJsonUtil.fromJson(lastLog.getRequestData(),...);callExternalApproval(param,reportId);}returnResult.success();}六、TransactionSynchronizationManager 原理这是 Spring 提供的事务同步管理器核心机制// Spring 事务提交流程简化publicvoidcommit(){// 1. 执行业务SQLdoCommit();// 2. 提交成功后调用所有注册的 synchronization.afterCommit()for(TransactionSynchronizationsync:synchronizations){sync.afterCommit();}// 3. 最终清理for(TransactionSynchronizationsync:synchronizations){sync.afterCompletion(STATUS_COMMITTED);}}publicvoidrollback(){// 回滚时不会调用 afterCommit()doRollback();// 只调用 afterCompletionfor(TransactionSynchronizationsync:synchronizations){sync.afterCompletion(STATUS_ROLLED_BACK);}}关键保证afterCommit()只在事务成功提交后才执行。如果事务回滚这个方法不会被调用。可用的生命周期钩子方法调用时机beforeCommit(boolean readOnly)事务提交前可以抛异常阻止提交afterCommit()事务提交成功后beforeCompletion()事务完成前提交或回滚都会调afterCompletion(int status)事务完成后status 区分提交/回滚七、操作日志表的设计意义每次与外部系统的交互都记录在日志表中┌────────────────────────────────────────────┐ │ report_submit_log │ ├────────────────────────────────────────────┤ │ id - 主键 │ │ report_id - 关联业务ID │ │ request_data - 发送数据JSON │ │ response_data - 返回数据JSON │ │ transmit_flag - T成功 / F失败 │ │ error_msg - 失败原因 │ │ create_time - 操作时间 │ └────────────────────────────────────────────┘作用可追溯— 出问题时可以看到发了什么、收到什么支持重试— 从日志中取出 request_data 重新调用问题定位— 是参数错了还是外部系统挂了数据恢复— 即使状态字段被意外修改日志表能还原真实历史八、状态流转图用户操作 MQ消费后 ┌───────────────────┐ ┌──────────────────────────────┐ │ │ │ │ │ DRAFT(草稿) │ │ 调用外部接口成功 │ │ │ │ │ → SUBMITTED(已提交) │ │ │ [提交] │ │ │ │ ▼ │ │ 调用外部接口失败 │ │ SUBMITTING │─MQ─→│ → SUBMIT_FAILED(提交失败) │ │ (提交中) │ │ │ │ │ │ 可重试 → 重新进入消费流程 │ └───────────────────┘ └──────────────────────────────┘提交中是一个过渡状态表示本地已经处理完毕但外部系统还没响应。这个中间状态的存在让系统能够阻止用户重复提交状态不是DRAFT了识别卡住的记录定时扫描超时的SUBMITTING正确显示进度前端可以展示处理中九、整体时序图用户 Controller Service(事务) DB MQ Consumer 外部系统 │ │ │ │ │ │ │ │──提交请求──→│ │ │ │ │ │ │ │──调用──→ │ │ │ │ │ │ │ │──校验查询──→│ │ │ │ │ │ │←─返回数据──│ │ │ │ │ │ │──更新状态──→│ │ │ │ │ │ │ (SUBMITTING)│ │ │ │ │ │ │──注册回调──→ (暂存) │ │ │ │ │ │──提交事务──→│ │ │ │ │ │ │ │ │ │ │ │ │ │──afterCommit触发──→ │ │ │ │ │ │ │ 发消息│ │ │ │←─返回成功──│←────────────│ │ │ │ │ │ │ │ │ │──消费──→ │ │ │ │ │ │ │ │──HTTP调用──→│ │ │ │ │ │ │←─返回结果──│ │ │ │ │←─更新状态(SUBMITTED)│ │ │ │ │ │←─保存日志──────────│ │十、这个模式的适用边界适用场景本地操作和外部调用需要最终一致性不需要强一致外部调用耗时较长500ms外部系统可能不稳定需要重试需要记录操作日志用于审计和排查不适用场景需要强一致性本地和外部必须同时成功或同时失败→ 考虑分布式事务Saga/TCC外部调用极快且稳定100ms→ 可以直接在事务内调用不需要异步用户必须等待外部结果 → 同步调用 超时处理潜在风险及应对风险应对方式事务提交成功但MQ发送失败定时任务扫描卡住的中间状态记录MQ消息重复消费Consumer 做幂等处理检查状态是否已变更外部系统长时间不可用重试次数限制 人工介入接口消息顺序问题同一业务ID的消息投递到同一队列分区十一、总结这个模式的核心要点事务内只做本地操作不调用外部系统通过afterCommit保证本地写成功后才触发下一步MQ 解耦异步调用外部系统Consumer 独立事务REQUIRES_NEW处理外部调用结果中间状态 操作日志支持追踪和重试定时补偿兜底 MQ 丢失的情况本质思想是将一个需要跨系统的复杂操作拆解为多个本地原子操作通过消息队列串联通过状态字段和日志表保证可追溯和可恢复。