
分布式事务几乎是 MQ 面试的必考点但也是很多候选人最容易含糊的地方。尤其是“事务消息”既要能说出 RocketMQ 半消息机制又要说清楚它解决的是最终一致性而不是强一致。面试官换一种问法比如“订单建好了库存没扣怎么办”“消息发送成功但本地事务回滚了怎么办”如果你只背了一套答案很容易卡住。这篇文章把 MQ 事务消息、分布式事务的核心要点整理成一套可复习的清单方案对比、工作原理、RocketMQ 代码示例、常见坑和面试回答思路。如果你正在准备面试或者工作中要在订单、库存、积分这类业务里做一致性设计这篇可以直接收藏。1. 分布式事务与 MQ 事务消息核心能力速览先看分布式事务常见方案的整体对比。面试时先说出这几种方案的区别再往事务消息上聚焦会显得有全局观。方案一致性类型对业务侵入性能影响典型组件适用场景2PC / XA强一致高高数据库、事务管理器跨库强一致但锁和协调者开销大TCC最终一致高中自研框架资金类、高频核心链路需要写Try/Confirm/Cancel本地消息表最终一致中低业务库 MQ通用异步场景兼容所有MQ事务消息最终一致中低RocketMQ通用异步场景推荐优先考虑最大努力通知最终一致低低MQ 定时任务通知类、对账类结果以查询为准关于 MQ 事务消息面试最核心的几个结论是事务消息的本质是“可靠消息”的一种实现解决的是消息发送与本地事务原子性问题。它不等于分布式事务的全部更不等于强一致。其中最关键的机制是半消息和事务状态回查。无论用哪种方案消费端都必须做幂等。2. 适用场景与使用边界事务消息不是万能的。先搞清楚它适合解决什么问题再决定要不要用它。典型适合场景下单成功后扣库存。支付成功后通知积分服务加积分。用户注册成功后发送优惠券。订单创建后触发物流、财务等异步流程。这些场景共同点是核心业务在本地事务中完成下游动作可以异步执行并且允许短时间内不一致最终通过消息重试达到一致。不适合场景账户扣款和余额更新要求实时强一致。两个服务必须同时成功或同时失败。下游失败会导致严重资金风险的场景。这类场景要优先考虑 TCC、本地事务或者干脆通过搜索系统状态而不是靠消息传递。使用边界也要说清楚事务消息不保证消费者只收到一次重试会带来重复投递。事务消息不保证下游一定能处理成功下游失败要靠消费重试和人工补偿。事务消息不保证消息延迟可控不能拿它做延迟队列。回查次数有上限超过后消息可能被丢弃必须有对账或告警兜底。从工程合规角度涉及订单、支付、库存等业务时不要在消息体中写入明文敏感信息发送前做好数据脱敏。生产链路要预留消息流水表方便核对和审计。3. MQ 事务消息工作原理半消息与回查机制先想一个问题为什么不能先发消息再执行本地事务如果先发消息本地事务失败消费者已经拿到消息去扣库存但订单没创建成功数据就错了。如果先执行本地事务再发消息消息发送失败时业务已经提交下游不知道需要扣库存。两种方式都做不到消息和本地事务一致。RocketMQ 的设计是先把消息发送到 Broker但这个消息处于半消息状态消费者不可见。业务本地事务执行完成后再通知 Broker 提交或回滚。完整流程如下Producer - Broker: 发送半消息 Producer - DB: 执行本地事务 Producer - Broker: COMMIT_MESSAGE / ROLLBACK_MESSAGE Broker - Producer: 如果状态未知发起事务状态回查 Producer - Broker: 返回本地事务真实状态 Broker - Consumer: 只有 COMMIT_MESSAGE 的消息才能被消费关键点是半消息半消息对消费者不可见。半消息也会持久化Broker 会单独管理。本地事务返回 COMMIT_MESSAGE消息才会变成普通消息。返回 ROLLBACK_MESSAGE消息会被删除。返回 UNKNOWNBroker 会按照配置定时向发送端发起回查。回查时发送端需要根据业务数据判断本地事务到底成功没有。比如订单服务收到回查就去查订单表订单存在则返回 COMMIT_MESSAGE不存在则返回 ROLLBACK_MESSAGE。RocketMQ 的回查次数和间隔受 Broker 端配置控制例如transactionCheckMax、transactionCheckInterval。回查超过上限仍然未知消息会被丢弃。生产环境一定要监控这个异常分支不能只依赖回查。4. 本地消息表 vs 事务消息两种可靠消息方案对比本地消息表是理解事务消息的最好前置知识也是面试常见追问点。本地消息表的核心思路是业务数据和消息状态放在同一个数据库事务里写。业务表插入一条订单消息表同时插入一条“待发送消息”这两个操作要么同时成功要么同时失败。然后由定时任务扫描消息表把待发送消息投递到 MQ投递成功后更新消息状态。消息表结构可以简化成这样CREATE TABLE local_message ( id BIGINT PRIMARY KEY AUTO_INCREMENT, biz_type VARCHAR(32) NOT NULL COMMENT 业务类型, biz_id VARCHAR(64) NOT NULL COMMENT 业务主键, payload TEXT NOT NULL COMMENT 消息内容, status TINYINT NOT NULL DEFAULT 0 COMMENT 0待发送 1已发送 2失败, retry_count INT NOT NULL DEFAULT 0 COMMENT 重试次数, create_time DATETIME NOT NULL, update_time DATETIME NOT NULL, UNIQUE KEY uk_biz_type_biz_id (biz_type, biz_id) ) COMMENT 本地消息表;定时任务发送逻辑伪代码public void scanAndSend() { ListLocalMessage messages messageMapper.selectPending(0, 100); for (LocalMessage message : messages) { try { SendResult result producer.send( new Message(ORDER_TOPIC, message.getPayload().getBytes()) ); if (result.getSendStatus() SendStatus.SEND_OK) { messageMapper.markSent(message.getId()); } } catch (Exception e) { messageMapper.incrementRetry(message.getId()); } } }本地消息表的优点是不依赖 MQ 的特定高级功能Kafka、RabbitMQ、RocketMQ 都能用。缺点也很明显业务代码里多维护一张表。定时任务需要处理发送失败、重试上限。消息重复投递概率更高消费端必须幂等。事务消息的不同在于把“本地事务 发送消息”的原子性交给 Broker 来协调。业务不用自己维护消息表而是通过监听器实现本地事务并通过回查机制兜底。对比表格如下对比项本地消息表事务消息对 MQ 的依赖任意 MQRocketMQ 等具备事务消息能力的 MQ业务侵入需要消息表 定时任务需要实现事务监听器事务一致性保障业务库 消息表同一事务半消息 本地事务 回查排查复杂度多一张表和一套任务需要理解半消息和回查机制通用性高中两种方案本质上都是“可靠消息 最终一致性”。面试时把本地消息表讲清楚再引出事务消息会显得对方案演进有理解。5. RocketMQ 事务消息落地实战这一节用一个经典场景订单创建成功后通知库存服务扣减库存。5.1 环境准备本地需要先准备 RocketMQ 运行环境JDK 8 及以上。RocketMQ 4.x 或 5.x。NameServer 和 Broker 能正常启动。如果使用管理后台dashboard需要单独启动 dashboard 项目。启动 NameServer 和 Broker 的通用命令如下# 启动 NameServer nohup sh bin/mqnamesrv # 启动 Broker并指定 NameServer 地址 nohup sh bin/mqbroker -n 127.0.0.1:9876 启动后可以用 mqadmin 检查集群状态sh bin/mqadmin clusterList -n 127.0.0.1:9876注意不同版本 RocketMQ 的启动脚本和参数可能略有差异生产环境不要直接裸跑需要配置账号鉴权、日志路径和存储路径。5.2 引入依赖以 Maven 为例dependency groupIdorg.apache.rocketmq/groupId artifactIdrocketmq-client/artifactId version4.9.4/version /dependency版本号以你本工程实际使用为准4.x 和 5.x 的事务消息 API 整体兼容。5.3 实现事务监听器事务监听器包含两个方法executeLocalTransaction执行本地事务并返回状态。checkLocalTransactionBroker 回查时调用确认本地事务结果。订单场景示例public class OrderTransactionListener implements TransactionListener { private final OrderMapper orderMapper; public OrderTransactionListener(OrderMapper orderMapper) { this.orderMapper orderMapper; } Override public LocalTransactionState executeLocalTransaction(Message msg, Object arg) { try { Order order (Order) arg; orderMapper.insert(order); return LocalTransactionState.COMMIT_MESSAGE; } catch (Exception e) { return LocalTransactionState.ROLLBACK_MESSAGE; } } Override public LocalTransactionState checkLocalTransaction(MessageExt msg) { String orderId msg.getKeys(); boolean exists orderMapper.existsByOrderId(orderId); return exists ? LocalTransactionState.COMMIT_MESSAGE : LocalTransactionState.ROLLBACK_MESSAGE; } }代码说明executeLocalTransaction中插入订单成功时返回 COMMIT_MESSAGE消息对消费者可见。失败时返回 ROLLBACK_MESSAGE半消息会被丢弃。如果本地事务执行期间进程重启或异常状态未知Broker 会回调checkLocalTransaction。回查逻辑以订单表数据为准这是最可靠的判断方式。5.4 发送事务消息使用TransactionMQProducer而不是普通DefaultMQProducerTransactionMQProducer producer new TransactionMQProducer( order_tx_producer_group ); producer.setNamesrvAddr(127.0.0.1:9876); producer.setTransactionListener(new OrderTransactionListener(orderMapper)); producer.start(); Order order buildOrder(); Message msg new Message( ORDER_STOCK_TOPIC, stock_deduct, order.getOrderId(), JSON.toJSONBytes(order) ); SendResult sendResult producer.sendMessageInTransaction(msg, order); System.out.println(sendResult.getSendStatus()); producer.shutdown();注意点sendMessageInTransaction的第二个参数会原样透传给executeLocalTransaction的arg。业务参数不要塞入消息体之外的敏感字段。发送事务消息后不能使用sendOneway否则无法确认半消息状态。5.5 消费端幂等消费RocketMQ 消费语义通常是 at-least-once网络抖动、重试都会导致重复消息。库存服务消费时一定要幂等。推荐做法在库存服务维护一张扣减流水表以消息 key 作为唯一键CREATE TABLE stock_deduct_log ( id BIGINT PRIMARY KEY AUTO_INCREMENT, order_id VARCHAR(64) NOT NULL, sku_id VARCHAR(64) NOT NULL, quantity INT NOT NULL, status TINYINT NOT NULL DEFAULT 1, create_time DATETIME NOT NULL, UNIQUE KEY uk_order_id (order_id) ) COMMENT 库存扣减流水表;消费逻辑示例consumer.registerMessageListener( (MessageListenerConcurrently) (msgs, context) - { for (MessageExt msg : msgs) { String orderId msg.getKeys(); if (deductLogMapper.existsByOrderId(orderId)) { continue; } // 先写流水再扣库存 deductLogMapper.insert(orderId); stockService.deduct(orderId); } return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; } );先写流水再执行扣减利用唯一索引挡住重复消息。即使扣减逻辑异常也可以根据流水表做补偿。6. Kafka 与 RabbitMQ 场景下的方案选择面试经常追问如果项目用的是 Kafka 或 RabbitMQ没有 RocketMQ 的事务消息怎么办6.1 Kafka 的事务 APIKafka 提供事务 API可以保证发送到多个分区的消息以及消费位移的提交是原子的。producer.initTransactions(); try { producer.beginTransaction(); producer.send(new ProducerRecord(order-topic, orderId, orderJson)); producer.sendOffsetsToTransaction(offsets, consumerGroupId); producer.commitTransaction(); } catch (Exception e) { producer.abortTransaction(); }这段代码解决了“处理消息 发送消息 提交消费位移”的原子性问题典型场景是 Kafka Streams 的恰好一次处理。但要注意Kafka 事务不能直接保证“业务数据库事务”和“发送消息”的原子性因为它无法感知外部数据库到底提交了没有。它和 RocketMQ 半消息模型不是一回事。Kafka 生态里做订单和库存一致性更常见的做法是业务数据库写订单同时写消息表和 Kafka 发送占位。定时任务或 CDC 把待发送记录投递到 Kafka。消费端做幂等。这本质上就是本地消息表模式。6.2 RabbitMQ 的事务与消息确认RabbitMQ 原生支持事务但性能较差channel.txSelect(); try { channel.basicPublish(exchange, routingKey, props, body); channel.txCommit(); } catch (Exception e) { channel.txRollback(); }事务模式下生产者发送是串行的吞吐下降明显。生产环境更推荐使用 Publisher Confirm 机制发送方开启 confirm。Broker 成功落盘后异步回调确认。发送方维护未确认消息失败后重发。消费端配合手动 ack 和幂等。这同样不能做到业务数据库和消息发送强原子仍需要在业务侧做补偿或使用本地消息表。面试回答时可以说没有哪种 MQ 是银弹重要的是根据业务一致性等级选择方案。如果需要全链路最终一致RocketMQ 事务消息、本地消息表、Kafka 事务加幂等三种方案都可行区别在于维护成本和延迟上限。7. 延迟消息与分布式事务的关系搜索热词里经常出现“mq 延迟消息队列”。很多候选人会把延迟消息和事务消息搞混这里单独说明。延迟消息解决的是“定时触发”问题比如下单 30 分钟未支付自动关单。支付完成后如果 10 秒内没收到对账回调触发主动查询。RocketMQ 支持延迟消息通过messageDelayLevel配置延迟等级。默认共有若干时间等级发送时指定等级Message msg new Message(ORDER_CLOSE_TOPIC, close_order, orderId, body); // 具体等级依赖 broker 的 messageDelayLevel 配置 msg.setDelayTimeLevel(16); producer.send(msg);注意不同 RocketMQ 版本的默认延迟等级可能不同生产环境以 Broker 配置为准。RocketMQ 的延迟消息是按等级设计的不能任意指定秒数。RabbitMQ 实现延迟队列通常有两种方式给消息设置 TTL消息过期后进入死信交换机。使用延迟消息插件rabbitmq-delayed-message-exchange。延迟消息不能替代事务消息事务消息保证的是消息与本地事务的一致性。延迟消息保证的是消息在指定时间后才被消费者看到。两者可以配合使用。例如订单模块订单创建成功后通过事务消息通知库存扣减。同时发送一条 30 分钟后的延迟消息用于判断订单是否超时未支付并主动关单。这样既能保证核心链路一致性又能做超时兜底。8. 发送端接口与批量投递的工程处理MQ 客户端通常提供三种发送方式同步发送、异步发送、单向发送。// 同步发送 SendResult sendResult producer.send(msg, 3000); // 异步发送 producer.send(msg, new SendCallback() { Override public void onSuccess(SendResult sendResult) { // 发送成功 } Override public void onException(Throwable e) { // 发送失败记录日志并重试 } }, 3000); // 单向发送不关心结果 producer.sendOneway(msg);事务消息必须使用sendMessageInTransaction不能使用单向发送。原因很简单半消息发送后需要依赖本地事务结果和回查来推进不能发送完就不管。批量投递场景需要注意RocketMQ 普通消息支持批量发送但批量消息有大小和条数限制。事务消息一般按单条发送逐条维护事务状态。批量扣库存时可以拆成多条普通消息加消费端幂等或者用同一事务消息携带列表由消费者逐条处理。发送失败重试建议同步发送失败可以捕获异常后重试但要注意幂等。消息 key 尽量用业务主键如订单 ID。重试超过阈值后记录日志并走补偿任务。不要把无限重试写在一个请求线程里容易拖垮应用。9. 性能、监控与资源消耗观察事务消息比普通消息多出的成本主要在三个方面半消息写入和回查消息带来的额外写盘。Broker 需要定时扫描半消息队列。发送方需要实现并响应回查接口。从监控角度需要重点观察Broker 磁盘使用率commitlog 增长是否异常。半消息 topic 中积压数量。回查次数如果回查过于频繁说明本地事务返回 UNKNOWN 太多。消费端消息堆积情况堆得越久最终一致性窗口越大。生产者发送耗时特别是事务消息发送sendMessageInTransaction的响应时间。常见的 mqadmin 命令可以检查集群状态和 topic 信息但不同版本命令参数不同sh bin/mqadmin clusterList -n 127.0.0.1:9876 sh bin/mqadmin topicList -n 127.0.0.1:9876 sh bin/mqadmin consumerProgress -g consumer_group -n 127.0.0.1:9876关于资源消耗不需要背具体数字。但面试时可以说明事务消息的额外开销与本地事务耗时、回查频率、日志保留时间相关。如果业务希望降低开销可以让executeLocalTransaction快速返回 COMMIT 或 ROLLBACK减少 UNKNOWN 状态从而减少回查。另外很多人在本地启动 RocketMQ 后管理后台打不开。这个问题大部分是环境问题排查顺序如下确认 NameServer 和 Broker 是否真的启动成功。确认 dashboard 项目里配置的namesrvAddr是否正确。确认管理后台端口是否被占用默认通常是 8080具体看项目配置文件。确认 9876 端口防火墙是否放行。10. 常见问题与排查方法问题现象可能原因排查方式解决方案事务消息一直不被消费半消息没有 COMMIT或回查逻辑异常查看半消息 topic、回查日志修正checkLocalTransaction保证返回真实状态本地事务成功但消息被回滚executeLocalTransaction 抛异常返回 ROLLBACK查看本地事务日志检查业务代码异常分支避免把外部调用失败当成业务失败消费者重复收到消息消费端重试导致重复投递查看消费日志和重投次数消费端幂等使用唯一索引或状态表库存扣减失败但消息已消费成功消费端没有重试机制或重试耗尽检查消费端异常日志配置消息重试失败进死信队列配合补偿任务管理后台无法进入NameServer/Broker 未启动或端口配置错误检查 9876、dashboard 端口和日志按排查顺序修复环境消息发送超时Broker 磁盘 IO 高或网络抖动查看 Broker 磁盘、发送耗时清理磁盘、扩容或优化发送方式回查次数过多本地事务耗时太长或返回 UNKNOWN 频繁查看回查日志控制本地事务耗时一次返回真实状态事务消息在 broker 重启后丢失半消息未刷盘或事务日志未同步检查 Broker 刷盘策略和主从配置生产环境开启适当刷盘策略并做集群部署面试时不需要把排查表背下来但至少要能说出事务消息出问题先看半消息状态再看回查日志最后看消费端幂等。11. 面试回答思路与最佳实践最后一个部分直接给出可用的面试回答框架。面试官问“订单创建后怎么保证库存一定扣减”推荐回答结构先分析场景。订单和库存允许短暂不一致可以接受最终一致不需要强一致。再选方案。如果 MQ 是 RocketMQ优先用事务消息。说原理。发送半消息执行本地事务插入订单返回 COMMIT消息才可见如果状态未知Broker 会回查。说实现。订单表写入和事务消息返回结果在一个事务语义内库存服务消费消息后先写扣减流水再扣库存保证幂等。说兜底。消费失败会重试重试失败进死信队列或对账任务极端情况靠人工补偿。面试官追问“如果消费者挂了一直消费不了怎么办”回答思路消息会重试RocketMQ 支持重试队列超过重试次数进入死信队列再配合定时任务扫描死信队列做补偿。更进一步的方案是业务层记录订单快照对账任务定期核对订单表和库存扣减流水。面试官追问“为什么不直接用 2PC”回答思路2PC 是强一致方案但跨服务场景锁资源、协调者单点、故障恢复复杂高并发下性能下降明显。订单扣库存这种场景最终一致已经够用没必要用强一致。最后给几条工程实践建议第一次接入事务消息先小额流量测试不要一上来全量切换。保留一套最小可运行配置方便排查环境问题。消息 key 使用业务主键日志里能串联完整链路。批处理任务要加日志、重试、失败告警。涉及订单、支付、库存的日志和消息体不打印敏感字段。消费端一律做幂等不信任 MQ 的至少一次语义。事务消息回查逻辑要查询业务真实状态而不是返回固定值。这套思路如果能在面试中顺畅讲出来并且能回答“为什么不用本地消息表”“为什么不用 Kafka 事务”基本就能把 MQ 事务消息和分布式事务这个考点拿稳。关键不是背代码而是把每个方案背后的取舍说清楚。