ARTICLE DETAIL

资讯详情

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

SpringCloud+RocketMQ事务消息:订单库存分布式事务落地

SpringCloud+RocketMQ事务消息:订单库存分布式事务落地 1. 从订单和库存说起分布式事务到底卡在哪做微服务的朋友大多都碰过这种场景用户下单订单服务写库成功库存服务扣减却失败或者反过来。单体架构里一个Transactional就能包住的东西拆成两个服务以后数据库本地事务就管不到对面了。我最早接手这类业务的时候leader给的方案很简单“你先扣库存扣完了再下单。”听起来没问题但真上线就翻车——扣库存成功、下单接口超时重试于是用户没订单库存却少了一件。要么就是先下单再扣库存高峰期并发一上来库存超卖得一塌糊涂。后来才明白跨服务的数据一致性本质上要靠分布式事务方案兜底而不是靠“谁先执行”来碰运气。这个场景里我最后选的是 SpringCloud RocketMQ 的分布式事务消息方案跑了大半年线上没有出过一笔对不上的账。这篇就把完整思路、底层原理、落地代码和踩过的坑一次性说清楚适合正在做订单、库存、支付这类对账敏感业务的团队参考。有人可能会问为什么不用 Seata 的 AT 模式或者 TCC这问题我后面会展开对比。简单先说结论RocketMQ 事务消息适合“本地事务 异步通知”这种天然能接受最终一致性的场景性能损耗小侵入性低而且 SpringCloud 生态里接入成本很可控。但在写代码之前我建议你先搞清楚一件事分布式事务方案没有银弹每一条路都在“一致性、可用性、性能”之间取舍。选 RocketMQ 事务消息是因为它最适合我们这个业务场景而不是因为它最“先进”。2. 方案选择题2PC、TCC、本地消息表、事务消息怎么挑2.1 四种主流方案的能力边界我先把做选型时对比过的四个方案摊开来说。2PC 两阶段提交数据库层面协调多个资源强一致但性能差、锁时间久而且协调者本身可能成为单点。微服务场景下跨库跨服务的 2PC 基本没人用在核心链路因为业务吞吐根本扛不住。我见过有人硬上 2PC压测一上 200 并发数据库锁等待直接把连接池打满最后只能连夜回滚。TCC 补偿事务Try-Confirm-Cancel 三段式每个服务都得写三套逻辑业务侵入极重。库存这种场景你得 Try 预扣、Confirm 确认扣减、Cancel 回补库存。代码量直接翻倍而且补偿逻辑里的脏数据处理很考验功底。适合对一致性要求极高的资金类业务但对普通订单库存来说成本和复杂度是过剩的。本地消息表业务表旁边建一个消息表事务里写业务和消息然后定时任务扫表发消息。这个方案很经典但坑也不少——定时扫描的延迟、消息表膨胀、发给 MQ 失败后的补偿逻辑全都要自己写。说白了它是把 MQ 的事务能力搬到了业务数据库里能用但不够优雅。RocketMQ 事务消息生产者先发一条“半消息”此时消费者不可见然后执行本地事务根据本地事务结果 commit 或 rollback 这条消息。如果本地事务执行过程中进程挂了RocketMQ 会定期回查事务状态自动兜底。这个机制刚好解决“本地事务和发消息不原子”的经典难题。2.2 为什么是 RocketMQ 而不是 Kafka 或 RabbitMQKafka 和 RabbitMQ 我都用过生态里也确实有类似方案但对比下来RocketMQ 的事务消息设计是最省心的。Kafka 的事务其实偏“流处理”场景它的事务 API 主要是保证“读-处理-写”这条链路的原子性跨系统跨服务的业务事务支持并不直接。RabbitMQ 有 publisher confirm 机制配合手动 ACK 能做可靠投递但“半消息 回查”这种专门为业务事务设计的原生能力还是 RocketMQ 做得最完整。还有一个实际感受RocketMQ 的事务回查机制是把补偿逻辑内置在 Broker 端的业务方只需要实现一个 callback 接口。而用 Kafka 或 RabbitMQ回查、补偿、重试基本都要自己在业务侧写一套调度逻辑维护成本完全不一样。所以选型结论如果你的业务能接受最终一致性且核心链路是订单、库存、积分这类异步解耦场景SpringCloud RocketMQ 事务消息是最稳的组合。3. RocketMQ 事务消息的底层逻辑半消息与回查机制3.1 半消息到底是什么先说“半消息”Half Message。这个名词听着玄乎其实逻辑很直白一条消息发到 Broker 后状态标记为“待确认”消费者正常情况下是拉不到这条消息的只有生产者后续告知 Broker“本地事务成功了”Broker 才会把消息翻转为可消费状态。整个过程分三步业务服务向 RocketMQ 发送半消息Broker 存储成功返回发送结果。业务服务收到发送成功的响应后执行本地数据库事务。本地事务执行完毕根据结果向 Broker 提交 commit 或 rollback 指令。这条链路里最关键的一点在于“发消息”和“执行本地事务”不再是一个先后的顺序关系而是像被一根线拴在了一起。哪怕本地事务执行到一半进程重启消息也不会丢——因为半消息还挂在 Broker 上等待后续的回查指令。我用一个生活化的类比来帮助你理解半消息就像你先递给餐厅一张“预约单”但厨房还不知道菜品是否真的要做等你确认了菜单、付了款那张预约单才会真正传到后厨开始做菜。如果付款失败预约单就自动作废菜品永远不会开始做。3.2 事务状态回查是怎么兜底的如果本地事务执行过程中服务直接宕机commit 或 rollback 指令永远发不出去那条半消息就会一直卡在 Broker 上。RocketMQ 的解决办法是事务回查。Broker 会按照默认的 60 秒间隔可配置去触发checkLocalTransaction回调业务方收到这个回调后去查一下自己本地事务的真实状态然后告诉 Broker 是 commit 还是 rollback。这里有一个非常重要的细节回查是由 Broker 主动发起的不是生产者自己定时检查。所以哪怕你的服务重启了只要恢复后还能被 Broker 连上半消息的最终命运就会由回查结果决定。我在代码里实现的checkLocalTransaction逻辑很简单根据事务消息里带的唯一业务订单号去查本地事务表中该订单的状态记录。查到状态为“已完成”就返回 COMMIT查不到或者标记为“已回滚”就返回 ROLLBACK边界情况则返回 UNKNOW 让 Broker 下次再查。注意回查次数是有限的默认最多查 15 次。超过次数仍然返回 UNKNOWBroker 会把消息丢进死信队列。这个限制逼着你在设计本地事务表的时候必须保证每条记录的状态是可追踪、可判定的这是事务消息方案能不能“稳”的关键。4. SpringCloud RocketMQ 工程落地从 pom 依赖到核心代码4.1 环境准备与基础依赖先交代我的工程环境SpringCloud版本用的 Hoxton.SR12SpringBoot 2.3.xRocketMQ 服务端 4.9.4RocketMQ 客户端用 4.9.4 保持版本对齐。Nacos 做注册中心和配置中心。如果你用更新的 SpringCloud Alibaba 版本记得核对依赖兼容性客户端版本和服务端版本建议大版本保持一致跨大版本很容易踩协议兼容的坑。服务端部署方面我在 CentOS 7 上装 RocketMQ 的时候踩过几个坑这里先提醒一下RocketMQ 默认的 JVM 内存配置是 4G小机器一定要改runserver.sh和runbroker.sh里的-Xms、-Xmx参数不然启动直接 OOM。另外如果你的 Broker 部署在云服务器上必须手动在broker.conf里配置brokerIP1否则客户端连上的是内网 IP生产消费者都会报连接异常。Maven 依赖只需要引入 rocketmq-spring-boot-starter 即可核心依赖如下dependency groupIdorg.apache.rocketmq/groupId artifactIdrocketmq-spring-boot-starter/artifactId version2.2.3/version /dependency这个 starter 封装了 RocketMQTemplate可以让 Spring Boot 工程像用 JmsTemplate 一样方便地发消息和注册监听器。配置文件里主要是 name-server 地址和生产组信息rocketmq: name-server: 192.168.1.100:9876 producer: group: order-producer-group send-message-timeout: 3000生产组名称必须和TransactionListener实现类挂钩这个细节后面会讲到先记住。4.2 订单服务本地事务与事务消息的原子绑定订单服务的核心逻辑是创建订单记录同时向 RocketMQ 发送事务消息通知库存服务扣减库存。关键在于订单写库和事务消息的发送必须由同一个TransactionListener来协调。我先把发消息的入口代码贴出来Service public class OrderService { Autowired private RocketMQTemplate rocketMQTemplate; Autowired private OrderMapper orderMapper; public void createOrder(OrderDTO orderDTO) { // 构造半消息消息体里带上订单号和明细 MessageString message MessageBuilder .withPayload(JSON.toJSONString(orderDTO)) .setHeader(orderId, orderDTO.getOrderId()) .build(); // 发送事务消息注意这里不是普通 send而是 sendMessageInTransaction rocketMQTemplate.sendMessageInTransaction( order-stock-topic, message, orderDTO ); } }sendMessageInTransaction是事务消息的核心入口。第三个参数会原样传给TransactionListener的executeLocalTransaction方法用来在本地事务里拿到上下文。这里不能直接使用rocketMQTemplate.convertAndSend因为普通 send 只有投递语义没有“本地事务确认”这个环节也就失去了分布式事务的意义。接着是实现TransactionListener接口这个接口是整个方案的心脏Component RocketMQTransactionListener(txProducerGroup order-producer-group) public class OrderTransactionListener implements TransactionListener { Autowired private OrderMapper orderMapper; Override public LocalTransactionState executeLocalTransaction(Message message, Object arg) { // arg 就是 sendMessageInTransaction 传入的 orderDTO OrderDTO orderDTO (OrderDTO) arg; try { // 本地事务插入订单主表和订单明细表 orderMapper.insertOrder(orderDTO); // 这里要实时返回 COMMIT告诉 Broker 可以放行消息 return LocalTransactionState.COMMIT_MESSAGE; } catch (Exception e) { // 本地事务失败返回 ROLLBACKBroker 会丢弃半消息 return LocalTransactionState.ROLLBACK_MESSAGE; } } Override public LocalTransactionState checkLocalTransaction(Message message) { String orderId message.getProperty(orderId); // 查本地事务表判断订单是否已经创建完成 OrderTransactionRecord record orderMapper.selectTransactionRecord(orderId); if (record ! null SUCCESS.equals(record.getStatus())) { return LocalTransactionState.COMMIT_MESSAGE; } if (record ! null FAIL.equals(record.getStatus())) { return LocalTransactionState.ROLLBACK_MESSAGE; } // 状态未知让 Broker 下次再查 return LocalTransactionState.UNKNOW; } }注意RocketMQTransactionListener注解里的txProducerGroup必须和配置里的producer.group一致否则事务回查不会触达你这个监听器。这个错不好排查因为消息能正常发出去消费者也能收到但“回查”永远不生效故障时你就发现消息状态没人管了。至于executeLocalTransaction里的操作有人会问“你这个本地事务方法直接执行了 insert那数据库事务呢”是的这块要配合DataSourceTransactionManager使用把订单表操作放在一个显式事务里执行或者用编程式事务。我实际项目里用了一段简单的TransactionTemplate包裹保证订单表插入的原子性。4.3 本地事务表回查机制的信息底座上面代码里提到了OrderTransactionRecord这可不是我写着玩的。RocketMQ 的回查机制本质上是 Broker 在问你“你那边的本地事务到底成没成”你得有一个可靠的信息源来回答它。我采用的方案是加一张事务记录表核心字段就五个tx_id消息里的唯一事务ID、order_id、statusSUCCESS/FAIL、create_time、update_time。事务记录和订单主表在同一次本地数据库事务里写入。这样回查时只需要根据 orderId 查这张表即可。有了这层表结构回查逻辑就干净了订单存在且状态 SUCCESS就 COMMIT订单创建失败或回滚标记就 ROLLBACK查不到任何记录可能是本地事务还没提交完成返回 UNKNOW等待下一轮回查。这里有个设计细节值得留意回查返回 UNKNOW 是有代价的——Broker 会暂停这条消息的消费释放并且占用回查配额。所以不要把回查做成“无脑 UNKNOW”要尽量在最大努力内判定清楚状态。我的做法是第一次回查查不到记录时等 3 秒再查一次还查不到才返回 UNKNOW这样能显著减少无效回查次数避免消息卡在中间状态过久。4.4 库存服务消息消费与幂等扣减消息从 Broker 放行后库存服务通过 RocketMQ 的监听器消费消息。这里有个坑我必须多说几句RocketMQ 事务消息只保证事务的最终一致性不保证消息只被消费一次。它实际投递语义是“至少一次”所以消费端必须做幂等。库存服务的消费逻辑如下Component public class StockConsumer { Autowired private StockService stockService; RocketMQMessageListener( topic order-stock-topic, consumerGroup stock-consumer-group ) public void onMessage(StockDeductMessage message) { // 幂等校验基于订单号商品SKU做去重 if (stockService.isHandled(message.getOrderId(), message.getSkuId())) { return; } // 加分布式锁防止并发重复扣减 boolean locked stockService.tryLock(message.getOrderId()); if (!locked) { // 锁没拿到可能是并发竞争尝试延迟重试 throw new RuntimeException(lock failed, retry later); } try { stockService.deductStock(message.getSkuId(), message.getQuantity()); stockService.markHandled(message.getOrderId(), message.getSkuId()); } finally { stockService.unlock(message.getOrderId()); } } }幂等落地的策略我觉得最靠谱的是两层结合第一层用 Redis 分布式锁挡住并发第二层在数据库里建一张stock_deduct_log表用订单号 SKU 作为唯一键。即使消费端重复执行唯一键冲突会让第二次插入直接失败不会再扣一次库存。这个方案上线后我做过故障演练手动重启库存服务模拟消息在消费过程中进程被杀然后重新拉起服务。RocketMQ 的消费进度没有 ACK消息会重新投递但因为幂等键的存在最终没有出现一笔重复扣减。4.5 客户端超时与发送失败的兜底事务消息的发送本身也可能失败比如 Broker 暂时不可用、网络抖动。sendMessageInTransaction抛异常时本地事务还没有执行此时可以直接向上抛错让调用方感知到“下单失败”不需要额外补偿。但有一个边界情况要注意消息发送成功但是返回响应在网络中超时了客户端会误判为失败而抛出异常。这时订单服务如果回滚本地事务RocketMQ 那边却已经收到了半消息回查机制会再次触发。由于本地事务表里没有对应的订单记录回查会返回 UNKNOW 直到超过次数消息进入死信队列。为了解决这个“假超时”问题我在发送消息前先往本地事务表插入一条状态为“PENDING”的记录。超时异常捕获后不立即回滚事务而是后续事务提交时把状态更新为 SUCCESS。回查时就能查到 PENDING 或 SUCCESSBroker 自然会把消息 commit 掉。这个处理方式说白了就是以本地事务表为准不以发送结果为准。分布式系统的很多坑根源都是“把网络响应当成事务结果”想通这一点设计上就能少踩很多坑。5. 实战中踩过的坑与排查体验5.1 事务回查不生效第一次接入时我在本地测试正常但线上部署后回查一直没有触发。排查了很久最后发现是RocketMQTransactionListener注解里的txProducerGroup和生产者的group不一致。我明明在配置文件里写了rocketmq.producer.group: order-producer-group但注解里写成了order-producer-group-test。RocketMQ 是依据生产者组来匹配回查监听器的对不上就静默地不回调连日志都没有。这类配置问题最难查因为表面看起来数据都能正常流转只有在服务宕机、需要回查兜底的时候才暴露。建议上线前做一次“宕机演练”发送一条半消息然后立即 kill 掉应用进程观察 Broker 回查触发后消息是否被正确 commit 或 rollback。5.2 消息消费重复导致库存扣两次上线后的第一个线上事故不是消息丢失而是消息重复消费。我的业务量不算大但库存表依然在一个晚上出现了两笔重复扣减。原因是消费端逻辑在“扣减库存”和“标记幂等”之间隔着一个 Redis 操作Redis 抖动导致标记失败了但库存已经扣了。后来我改成“先标记幂等再扣库存”并且把标记和扣减放在同一个本地数据库事务中通过唯一键来保证原子性。Redis 只是并发控制的辅助手段不再是幂等的最终依据。真实的生产环境任何中间件都有抖动风险唯一能信得过的就是数据库的约束。5.3 回查频率过高拖垮数据库默认情况下 RocketMQ 回查间隔是 60 秒看起来不高但如果积压的半消息几百上千条Broker 的回查线程会集中打过来数据库会瞬时负载飙高。我遇到过凌晨 MQ 积压早上业务恢复消费时回查风暴把订单库的 CPU 打到 100%。两个优化点第一把回查间隔调大到 120 秒给本地事务更多的完成时间也降低数据库压力第二回查方法里直接走“内存缓存 数据库兜底”的二级策略刚提交的事务能在缓存里查到就直接返回 COMMIT不再打数据库。这个策略把订单库在回查场景下的 QPS 降了一个量级。5.4 消费者组同一个 Topic 多监听器冲突如果同一个 Topic 被多个消费者组监听消息会被广播给所有组没问题。但如果一个消费者组内的多个RocketMQMessageListener监听了同一个 TopicRocketMQ 会把消息随机分配给其中一个实例导致逻辑错乱。我踩过这个坑之后强制规定一个消费者组对应一个 Topic 和一套消费逻辑分组按业务模块独立不要共用。5.5 消息顺序问题事务消息默认不是严格有序的。订单创建和取消这两条消息如果走了不同的队列消费端可能先收到取消再收到创建。处理这类逻辑我通常在消息体里带一个业务时间戳消费时做版本比对只处理“最新状态”的操作。这是比较轻量的方案如果对顺序要求更强需要把同一订单的消息路由到同一个消息队列RocketMQ 的MessageQueueSelector可以做到但会牺牲一定的并发度。6. 一套可以直接抄作业的检查清单分享几个我在验收这套方案时会逐项过的点你可以拿来当自检清单一致性验证本地事务表和业务主表是否在同一数据库事务中写入回查判定是否能覆盖所有分支成功、失败、未知是否对“消息发送超时但 Broker 已收到”的场景做了兜底幂等验证消费端是否依赖数据库唯一键而不是 Redis 来保证幂等重复消费时是否不会产生副作用重复扣库存、重复加积分可靠性验证消费者组是否开启了重试机制重试次数和间隔是否合理死信队列是否有告警通知和处理流程服务重启后未完成的事务消息是否能被回查机制继续处理性能验证回查频率是否过高是否有缓存策略保护数据库半消息积压量是否有监控指标每一轮的故障演练都建议做一遍停掉 Broker、杀掉生产者、杀掉消费者、断网重连。把这些场景都过一遍这套方案才算真正在你手里“稳的一批”。7. 最后分享一点个人体会用 RocketMQ 事务消息做分布式事务最大的体会是“设计上做减法”。很多团队一上来就想上 TCC、上 Seata 全局事务结果代码复杂度暴涨排障成本成倍增加。事务消息的核心思想是“先把本地事务做完再通知别人”这种最终一致的思路正好匹配互联网业务的大部分场景。我个人建议是先在低风险业务上跑通这套方案积累经验后再推广到核心链路。不要一上来就拿订单库存这种最敏感的业务练手。另外事务消息不是万能的如果你的业务需要强一致比如账户余额和交易的实时对账事务消息的异步窗口可能满足不了要求那时候再考虑 TCC 或本地消息表也不迟。这套 SpringCloud RocketMQ 的组合我前后用了两个项目累计跑了一年多从几千单到几十万单都扛住了。只要你把半消息机制理解透、本地事务表设计好、消费端幂等做到位它确实配得上“稳的一批”这四个字。
返回列表