ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

线上问题处理记录:事务未提交 + 异步消息消费:竞态问题全解析

线上问题处理记录:事务未提交 + 异步消息消费:竞态问题全解析 线上问题处理记录事务未提交 异步消息消费竞态问题全解析一、问题现象调拨单创建后应自动生成待发货单但有一部分调拨单一直停留在待发货状态却始终没有对应的待发货单记录。二、问题发生流程时间轴 → Thread-A (业务线程) Thread-B (MQ 消费线程) ──────────────────── ────────────────────── BEGIN TRANSACTION ① INSERT xxx_master (调拨单) ② 发送 MQ 消息到 xxx-wait-delivery 队列 ③ 收到消息开始消费 ④ SELECT xxx_master WHERE order_code? ⑤ 查不到数据 → 直接 return跳过 COMMIT TRANSACTION此时数据才可见核心矛盾消息在事务提交之前就被发送并消费了。消费者查询数据库时由于事务隔离性读取不到尚未提交的数据。注博客https://blog.csdn.net/badao_liumang_qizhi三、根本原因分析3.1 事务与消息的时序错位在同一个事务方法中写入数据库此时尚在事务内其他线程不可见发送 MQ 消息消息立即投递到 Broker方法结束事务提交数据才对外可见MQ 投递几乎瞬时而事务提交取决于数据库 I/O。如果消费者响应足够快就会在数据可见之前去读取造成幻读。3.2 数据库隔离级别的影响MySQL 默认隔离级别为REPEATABLE READ消费者开启新事务后读到的是该事务开始时刻的快照即便生产者事务稍后提交了消费者的快照中也不包含这条新记录3.3 MQ 堆积时矛盾进一步放大正常情况下消费者可能比事务提交慢几十毫秒表面上不会出问题。但当队列堆积 → 消费追赶 → 突然并发大量消费时时序重合概率显著增大。四、涉及的技术点技术点说明本地事务 ACID事务提交前写入对其他会话不可见SpringTransactional方法级事务边界方法内所有操作在同一事务MQ 异步解耦消息发送后即进入 Broker与发送方事务无关事务同步回调Spring 提供TransactionSynchronization在提交后执行动作最终一致性分布式场景下无法同时保证强一致需选择合适的补偿策略幂等消费消费者需要处理消息重复投递的场景五、正确的解决方案方案一事务提交后再发送消息推荐importorg.springframework.transaction.support.TransactionSynchronization;importorg.springframework.transaction.support.TransactionSynchronizationManager;TransactionalpublicvoidcreateOrderAndNotify(Orderorder){// 1. 数据库写入事务内orderRepository.save(order);// 2. 注册事务提交后回调TransactionSynchronizationManager.registerSynchronization(newTransactionSynchronization(){OverridepublicvoidafterCommit(){// 事务已提交数据对外可见此时发送 MQmqSender.send(order-queue,order.getId());}});}优点消费者消费时数据一定已存在。缺点如果发送 MQ 失败消息丢失需要配合定时补偿任务。方案二消费者增加前置校验 重试RabbitListener(queuesorder-queue)publicvoidconsume(StringorderId){OrderorderorderRepository.findById(orderId);if(ordernull){// 数据尚未可见抛异常触发 MQ 重试thrownewRetryableException(订单不存在等待重试: orderId);}// 正常业务处理processOrder(order);}优点兼容各种时序异常不修改发送端。缺点依赖 MQ 重试机制若数据确实不存在已删除需设置最大重试次数避免死循环。方案三事务消息RocketMQ / 本地消息表// 伪代码本地消息表方案TransactionalpublicvoidcreateOrder(Orderorder){orderRepository.save(order);// 消息写入本地消息表与业务数据同一事务localMessageRepository.save(newLocalMessage(order-queue,order.getId()));}// 定时任务扫描本地消息表发送到 MQ 后标记已发送Scheduled(fixedDelay1000)publicvoidscanAndSend(){ListLocalMessagependinglocalMessageRepository.findPending();for(LocalMessagemsg:pending){mqSender.send(msg.getQueue(),msg.getPayload());msg.markSent();localMessageRepository.save(msg);}}优点保证消息与数据的强一致性。缺点实现复杂增加数据库写入和定时扫描开销。方案对比方案一致性保证实现复杂度适用场景事务提交后发送较强极小窗口丢消息低大多数业务消费者校验重试最终一致低允许短暂延迟本地消息表强一致高金融/核心链路六、避免踩坑清单✅ DO永远在事务提交后发送 MQ使用TransactionSynchronization.afterCommit()或 Spring 4.2 的TransactionalEventListener(phase AFTER_COMMIT)。消费者对关键数据做存在性校验查不到时触发可控重试延迟队列/指数退避。消费者保证幂等同一消息消费多次结果一致。设置最大重试次数和死信队列防止无限重试。MQ 堆积期间评估状态变更影响消费前校验数据当前状态是否仍然合法。❌ DON’T不要在Transactional方法内直接mqTemplate.send()。不要假设消费一定比事务提交慢——网络抖动、GC 暂停、队列堆积后追赶消费都会打破这个假设。不要把 MQ 消费逻辑写成查不到就忽略——数据可能只是还没提交不是真不存在。不要在没有幂等保证的情况下无限重试。七、类似踩坑案例案例 1支付回调 vs 订单创建现象第三方支付回调先于订单写入完成导致回调处理时找不到订单。原因调用第三方支付接口在事务内支付方异步回调比本地事务提交快。解决回调接收后先入临时表由定时任务匹配已存在的订单处理。案例 2主从延迟 异步查询现象写入主库后立即发送 MQ消费者从从库读取因主从延迟读不到。原因读写分离场景下从库有毫秒到秒级延迟。解决消费者强制走主库查询或在消息体中携带完整数据而非仅 ID。案例 3分布式事务 事件通知现象A 服务调用 B 服务写入数据后发事件C 服务消费事件去 B 服务查询偶发查不到。原因B 服务接口返回成功但事务尚未提交异步 Flush 或二阶段提交延迟。解决B 服务在事务提交后自行发布事件事件归属于数据拥有者。案例 4批量操作 长事务现象批量导入 1000 条记录事务耗时 5 秒期间部分 MQ 已消费完毕。原因大事务长时间不提交消息早已投递。解决拆分为小批次事务每批提交后再发送对应消息。八、通用示例代码Spring Boot RabbitMQ8.1 生产者事务提交后发送ServiceSlf4jpublicclassOrderService{privatefinalOrderRepositoryorderRepository;privatefinalRabbitTemplaterabbitTemplate;publicOrderService(OrderRepositoryorderRepository,RabbitTemplaterabbitTemplate){this.orderRepositoryorderRepository;this.rabbitTemplaterabbitTemplate;}TransactionalpublicOrdercreateOrder(CreateOrderRequestrequest){// 1. 业务数据持久化OrderordernewOrder();order.setOrderCode(request.getOrderCode());order.setStatus(OrderStatus.CREATED);order.setCreateTime(newDate());orderRepository.save(order);// 2. 事务提交后发送 MQTransactionSynchronizationManager.registerSynchronization(newTransactionSynchronization(){OverridepublicvoidafterCommit(){try{rabbitTemplate.convertAndSend(order-exchange,order.created,order.getId());log.info([MQ] 事务提交后发送消息成功, orderId{},order.getId());}catch(Exceptione){// 发送失败记录日志由补偿任务处理log.error([MQ] 消息发送失败, orderId{},order.getId(),e);}}});returnorder;}}8.2 消费者前置校验 延迟重试ComponentSlf4jpublicclassOrderConsumer{privatestaticfinalintMAX_RETRY5;privatefinalOrderRepositoryorderRepository;privatefinalRabbitTemplaterabbitTemplate;publicOrderConsumer(OrderRepositoryorderRepository,RabbitTemplaterabbitTemplate){this.orderRepositoryorderRepository;this.rabbitTemplaterabbitTemplate;}RabbitListener(queuesorder-process-queue)publicvoidhandleOrderCreated(Messagemessage){StringorderIdnewString(message.getBody());intretryCountgetRetryCount(message);// 1. 前置校验数据是否已可见OrderorderorderRepository.findById(Integer.valueOf(orderId)).orElse(null);if(ordernull){if(retryCountMAX_RETRY){log.error([消费] 达到最大重试次数仍查不到数据, orderId{}, 转入死信,orderId);// 不再重试消息将进入死信队列thrownewAmqpRejectAndDontRequeueException(数据不存在: orderId);}log.warn([消费] 数据尚未可见, orderId{}, 第{}次重试,orderId,retryCount1);// 投递到延迟队列等待后重新消费sendToDelayQueue(orderId,retryCount1);return;}// 2. 状态校验数据存在但可能已被取消if(order.getStatus()OrderStatus.CANCELLED){log.info([消费] 订单已取消, 跳过处理, orderId{},orderId);return;}// 3. 幂等校验防止重复消费if(order.getStatus()OrderStatus.PROCESSED){log.info([消费] 订单已处理, 跳过重复消费, orderId{},orderId);return;}// 4. 正常业务处理processOrder(order);log.info([消费] 订单处理完成, orderId{},orderId);}privatevoidsendToDelayQueue(StringorderId,intretryCount){rabbitTemplate.convertAndSend(order-exchange,order.delay,orderId,msg-{// 延迟 3 秒后重新投递需配合 RabbitMQ 延迟插件或 TTL 死信路由msg.getMessageProperties().setDelay(3000*retryCount);msg.getMessageProperties().setHeader(x-retry-count,retryCount);returnmsg;});}privateintgetRetryCount(Messagemessage){Objectcountmessage.getMessageProperties().getHeader(x-retry-count);returncountnull?0:(int)count;}privatevoidprocessOrder(Orderorder){// 实际业务逻辑order.setStatus(OrderStatus.PROCESSED);orderRepository.save(order);}}8.3 使用 TransactionalEventListener更简洁的写法// 事件定义publicclassOrderCreatedEvent{privatefinalIntegerorderId;publicOrderCreatedEvent(IntegerorderId){this.orderIdorderId;}publicIntegergetOrderId(){returnorderId;}}// 生产者ServiceTransactionalpublicclassOrderService{privatefinalOrderRepositoryorderRepository;privatefinalApplicationEventPublishereventPublisher;publicOrdercreateOrder(CreateOrderRequestrequest){OrderordernewOrder();order.setOrderCode(request.getOrderCode());order.setStatus(OrderStatus.CREATED);orderRepository.save(order);// 发布领域事件此时不会立即触发监听器eventPublisher.publishEvent(newOrderCreatedEvent(order.getId()));returnorder;}}// 事件监听器仅在事务提交后触发ComponentSlf4jpublicclassOrderEventListener{privatefinalRabbitTemplaterabbitTemplate;TransactionalEventListener(phaseTransactionPhase.AFTER_COMMIT)publicvoidonOrderCreated(OrderCreatedEventevent){rabbitTemplate.convertAndSend(order-exchange,order.created,event.getOrderId());log.info([Event] 事务提交后发送 MQ, orderId{},event.getOrderId());}}九、总结要素本次问题通用规律根因事务内发送 MQ消费先于提交跨系统操作无法共享事务边界表现消费者查不到数据直接跳过下游读不到上游写入修复消费端增加状态校验发送端移到事务后afterCommit 消费前校验双保险兜底运维接口批量补偿定时任务/人工补偿覆盖极端场景核心原则谁拥有数据谁负责在数据可见后通知外部消费者永远不要假设数据一定存在要有校验和重试机制。
返回列表