面试问 MQ 事务消息,最怕不是背不出概念,而是只能回答一句“RocketMQ 支持事务消息”。真正拉开差距的是:你能不能把“半消息、本地事务、回查机制、最终一致性、幂等消费”这一条链路讲清楚,能不能现场写出一个订单与库存的示例,能不能说出消息丢失、重复消费、延迟消息这些衍生问题怎么处理。这篇文章就把这些点一次性讲透。
文中会先给出 MQ 方案在分布式事务里的定位和选型速览,然后拆解事务消息的原理,再用 RocketMQ 事务消息、本地消息表、幂等消费三个代码示例演示完整设计,最后补上延迟消息补偿、管理后台排查、面试高频追问和使用建议。适合准备 MQ 相关岗位面试的开发者,也适合正在做订单、库存、支付、积分这类最终一致性方案的工程人员。
1. 核心能力速览
先说结论:MQ 不是用来实现强一致的,它是用来实现分布式系统最终一致性的。事务消息是 MQ 在“基础可靠传输”之上提供的一种保证,核心目标是“本地事务和消息发送要么都成功,要么都失败”。
不同 MQ 产品对事务消息的支持程度不一样,这里做一张快速对比表:
| 对比项 | RocketMQ | Kafka | RabbitMQ | Pulsar |
|---|---|---|---|---|
| 原生事务消息 | 支持 | 支持事务 API,用途偏向“精确一次消费” | 不支持原生事务消息,用本地消息表方案 | 支持事务 |
| 延迟消息 | 支持延迟等级/定时消息 | 不原生支持,需自研或用时间轮 | 通过延迟消息插件支持 | 支持定时消息 |
| 实现复杂度 | 中等,消息已封装 | 较高,事务 API 理解成本高 | 低,但事务场景要改架构 | 中等 |
| 典型场景 | 订单、交易、支付、库存等电商链路 | 大数据链路、日志、埋点、流处理 | 内部系统异步通知、事件订阅 | 云原生、跨地域复制 |
| 客户端生态 | Java 为主 | 多语言完善 | 多语言完善 | 多语言完善 |
面试时说 MQ 事务消息,一般默认指 RocketMQ 的事务消息。如果是 Kafka 或 RabbitMQ,需要单独说明替代方案,比如 Kafka 配合幂等 Producer 和事务 API,RabbitMQ 配合 Publisher Confirm 加本地消息表。
2. 前置知识:分布式事务的几种典型方案
在聊 MQ 事务消息之前,必须先铺垫分布式事务的整体框架。否则面试官一追问“为什么不用 2PC”就会卡住。
常见分布式事务方案有五类:
- 2PC(两阶段提交):强一致,但协调者单点、阻塞时间长,性能差。
- TCC(Try-Confirm-Cancel):业务侵入强,每个操作都要写三个接口,适合资金类强约束场景。
- 本地消息表:基于数据库事务写业务表和消息表,再通过定时任务扫描发送。
- MQ 事务消息:消息中间件替我们实现“本地事务和消息发送的原子性”。
- Sagas 长事务:通过一系列本地事务和补偿事务完成,适合流程长、允许补偿的场景。
MQ 的定位,是让“本地数据库事务”和“异步消息通知”保持最终一致,不追求同步强一致。典型的落地场景就是订单和库存:用户下单时,订单系统写入订单表,同时发一条 MQ 消息通知库存系统扣减库存。如果先写订单再发消息,可能消息发失败;如果先发消息再写订单,可能出现订单没创建成功但库存已经扣了。事务消息就是为了解决这两个操作“要么都成功,要么都失败”的问题。
3. 什么是 MQ 事务消息
以 RocketMQ 为例,事务消息的流程分三个阶段。
第一阶段,生产者发送一条“半消息”(Half Message)到 Broker。半消息和普通消息的区别是:它对消费者不可见,处于暂存状态。
第二阶段,生产者执行本地事务,也就是写订单、锁库存这些真实业务操作。
第三阶段,生产者根据本地事务的执行结果,向 Broker 提交 commit 或 rollback。commit 后半消息变成可见消息,消费者才能拉取;rollback 后半消息被删除,消费者永远看不到。
这里面有一个关键问题:如果本地事务执行完了,但发送 commit/rollback 时进程挂了怎么办?所以 Broker 有一个“事务回查”机制。Broker 会定期反问生产者“你这笔本地事务到底成没成”,生产者的 TransactionListener 中要实现 checkLocalTransaction 方法,根据业务数据判断应该 commit 还是 rollback。
整体链路可以这样理解:
- 发送半消息。
- 执行本地事务。
- 提交或回滚消息。
- Broker 回查兜底。
- 消息可见后,消费者消费。
- 消费者侧做幂等,保证不重复扣库存。
这套机制保证了:业务数据库里的订单记录,和 MQ 里的最终可见消息,是一致的。不会出现订单成功但消息没发出去,也不会出现订单失败但消息对消费者可见。
4. 订单与库存分布式事务设计
订单和库存是最经典的分布式事务场景。下面画一条完整时序路径。
用户下单后,订单服务做两件事:
- 在订单库写入订单数据。
- 发送事务消息,通知库存服务扣减库存。
库存服务收到消息后,执行库存扣减,并返回结果。如果库存不足,则抛出业务异常,触发告警或人工补偿。
这里要注意,消费者即使收到消息,也可能执行失败。比如库存扣减超时、数据库锁等待、服务重启等。所以消费端必须做好“失败重试 + 幂等”:
- 消费失败,MQ 会按重试队列重新投递。
- 重复消费时,通过消息唯一 ID 或业务唯一键防止重复扣减。
如果库存始终扣减失败,需要有一个补偿机制。常见的做法是使用延迟消息重试,先延迟 10 秒,再延迟 30 秒,逐渐增加重试间隔。重试几次仍失败,就落到人工处理表,由运营介入。
这套设计不是强一致,而是最终一致。用户下完单可能短暂看到库存没有实时扣减,但经过消息消费后,最终库存会被正确扣减。这就是 MQ 方案和 2PC 方案的本质区别。
5. 代码实践:RocketMQ 事务消息
面试手写代码时,不需要写完整的项目工程,但核心 TransactionListener 一定要能默写出来。
先看生产者发送事务消息的代码。以 RocketMQTemplate 为例:
@Resource private RocketMQTemplate rocketMQTemplate; public void createOrder(OrderDO order) { String topic = "order-tx-topic"; Message<String> message = MessageBuilder.withPayload(JSON.toJSONString(order)) .setHeader("orderId", order.getOrderId()) .build(); TransactionSendResult result = rocketMQTemplate.sendMessageInTransaction( topic, message, order ); if (LocalTransactionState.COMMIT_MESSAGE.equals(result.getLocalTransactionState())) { log.info("订单消息已提交, orderId = {}", order.getOrderId()); } else { log.warn("订单消息未提交, orderId = {}", order.getOrderId()); } }第三参数 arg 是透传对象,在 TransactionListener 的 executeLocalTransaction 里可以直接拿到。这就是你传订单实体进去的原因。
接着实现 TransactionListener:
@Component public class OrderTransactionListener implements TransactionListener { @Resource private OrderMapper orderMapper; @Resource private InventoryMapper inventoryMapper; @Override public LocalTransactionState executeLocalTransaction(Message msg, Object arg) { try { // 这里才是真正的本地事务 OrderDO order = (OrderDO) arg; orderMapper.insert(order); inventoryMapper.preReduce(order.getSkuId(), order.getCount()); return LocalTransactionState.COMMIT_MESSAGE; } catch (Exception e) { log.error("本地事务执行失败", e); // 返回 UNKNOW,等待 Broker 回查 return LocalTransactionState.UNKNOW; } } @Override public LocalTransactionState checkLocalTransaction(MessageExt msg) { // Broker 回查时,根据订单是否已经写入来判断事务结果 String orderId = msg.getKeys(); OrderDO order = orderMapper.selectByOrderId(orderId); if (order != null) { return LocalTransactionState.COMMIT_MESSAGE; } return LocalTransactionState.ROLLBACK_MESSAGE; } }这个代码里有两个关键细节。
第一,executeLocalTransaction 中如果业务操作抛出异常,不能直接返回 ROLLBACK_MESSAGE,而是要返回 UNKNOW。因为异常可能是临时网络抖动或数据库连接失败,直接回滚半消息会让“本地事务是否成功”的最终判断丢失。返回 UNKNOW 后 Broker 会回查,让系统自己核实订单到底成没成。
第二,checkLocalTransaction 的查询一定要快。因为 Broker 回查是有频率和次数限制的,如果每次回查都查库超时,消息会一直处于中间状态,最终会被丢弃或进入死信队列。查询建议走主键或唯一索引。
6. 本地消息表方案:适合不支持事务消息的 MQ
如果你们公司用的是 RabbitMQ,或者用的 MQ 版本不支持事务消息,不要慌,还有经典方案:本地消息表。
核心思路是这样的:在一次数据库事务里,同时写业务表和消息表。业务提交成功,消息表一定也有一条对应记录。然后由定时任务扫描消息表,把未发送的消息发给 MQ。消费者消费完成后,再通知消息表更新状态。
先建一张消息表:
CREATE TABLE mq_message_record ( id BIGINT PRIMARY KEY AUTO_INCREMENT, msg_id VARCHAR(64) NOT NULL UNIQUE, biz_key VARCHAR(64) NOT NULL, topic VARCHAR(128) NOT NULL, payload VARCHAR(2048) NOT NULL, status TINYINT NOT NULL DEFAULT 0 COMMENT '0-待发送 1-已发送 2-已确认', retry_count INT NOT NULL DEFAULT 0, next_retry_time DATETIME NOT NULL, create_time DATETIME NOT NULL, update_time DATETIME NOT NULL, KEY idx_status_retry (status, next_retry_time) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;业务侧事务长这样:
@Transactional(rollbackFor = Exception.class) public void createOrder(OrderDO order) { // 1. 写业务数据 orderMapper.insert(order); // 2. 写消息表 MqMessageRecord record = MqMessageRecord.builder() .msgId(UUID.randomUUID().toString().replace("-", "")) .bizKey(order.getOrderId()) .topic("order-tx-topic") .payload(JSON.toJSONString(order)) .status(0) .build(); mqMessageRecordMapper.insert(record); }定时任务扫描待发送消息:
@Scheduled(fixedDelay = 5000) public void sendPendingMessage() { List<MqMessageRecord> pendingList = mqMessageRecordMapper.selectPending(100); for (MqMessageRecord record : pendingList) { try { SendResult sendResult = mqProducer.send(record.getTopic(), record.getPayload(), record.getMsgId()); if (sendResult.getSendStatus() == SendStatus.SEND_OK) { mqMessageRecordMapper.markSending(record.getId()); } } catch (Exception e) { log.error("消息发送失败, msgId = {}", record.getMsgId(), e); mqMessageRecordMapper.markRetry(record.getId()); } } }这个方案的优点是:对 MQ 本身零依赖,RabbitMQ、Kafka 都能用。缺点是:消息表会随着业务增长变大,需要定期清理;定时任务扫描会带来秒级延迟;消息发送语义仍是“至少一次”,消费者必须幂等。
面试如果被问到“你们为什么不用 RocketMQ 事务消息”,你可以说:已有 MQ 基础设施是 RabbitMQ,为了避免引入新的中间件,使用本地消息表;如果公司已经在用 RocketMQ,优先考虑原生产品能力。
7. 幂等消费:分布式事务一致性的最后一道防线
无论用哪种方案,MQ 的消息都可能重复。网络超时重发、Broker 重试、消费者重启都会导致重复消息。所以分布式事务消息方案永远要配幂等。
判断一个消费端是否幂等,标准是:同一个业务请求执行一次和执行多次,结果一致。
库存扣减这种操作最怕重复。用户只下单一次,库存却被扣两次,这属于严重事故。因此消费端必须加唯一约束。
一种做法是使用数据库唯一键。消息表的 msg_id 唯一,消费端先尝试插入消费记录,如果唯一键冲突,说明处理过了,直接跳过:
@Transactional public void onInventoryMessage(InventoryMessage message) { try { consumeRecordMapper.insert(ConsumeRecord.builder() .msgId(message.getMessageId()) .bizKey(message.getOrderId()) .build()); } catch (DuplicateKeyException e) { log.info("重复消息已忽略, msgId = {}", message.getMessageId()); return; } inventoryMapper.reduce(message.getSkuId(), message.getCount()); }另一种做法是用 Redis 的 setIfAbsent 做前置判重,适合对性能要求高的场景:
public void onInventoryMessage(InventoryMessage message) { String key = "idem:" + message.getMessageId(); Boolean success = redisTemplate.opsForValue().setIfAbsent(key, "1", Duration.ofMinutes(5)); if (Boolean.TRUE.equals(success)) { inventoryMapper.reduce(message.getSkuId(), message.getCount()); } else { log.warn("重复消息已被过滤, msgId = {}", message.getMessageId()); } }Redis 判重的优点是快,缺点是不能覆盖 Redis 超时和 Redis 宕机。所以核心资金链路建议“数据库唯一键 + Redis 前置过滤”双保险。
注意:幂等不能只依赖 MQ 自带的 msgId。同一个业务操作在重试时可能生成新的 msgId,必须使用业务唯一键,例如 orderId、paymentId、bizKey。这也是面试官最常挖的细节。
8. 延迟消息与重试补偿设计
最新搜索热词里反复出现“mq 延迟消息队列”,这个话题和分布式事务关系非常紧密。重试补偿就是靠延迟消息实现的。
以 RocketMQ 为例,延迟消息通过设置 delayLevel 实现:
Message message = new Message("inventory-compensate-topic", "orderId", orderId.getBytes()); message.setDelayTimeLevel(3); producer.send(message);RocketMQ 有固定的延迟等级,例如 1s、5s、10s、30s、1m、2m 等。不同版本等级数字不一样,使用时需要查对应版本文档。
在订单与库存场景,延迟消息通常这样用:库存扣减失败后,不立即无限重试,而是发送一条延迟消息到补偿队列,延迟 10 秒后再消费。如果还是失败,再延迟 30 秒。达到最大重试次数后,写入人工补偿表,同时给运营发告警。
@Scheduled(cron = "0 */1 * * * ?") public void compensateInventory() { List<InventoryAdjustRecord> failedList = inventoryAdjustRecordMapper.selectFailed(100); for (InventoryAdjustRecord record : failedList) { InventoryCompensateMessage msg = new InventoryCompensateMessage(); msg.setOrderId(record.getOrderId()); msg.setSkuId(record.getSkuId()); msg.setCount(record.getCount()); msg.setRetryCount(record.getRetryCount()); rocketMQTemplate.syncSend("inventory-compensate-topic", msg, 3000); } }如果使用的 MQ 没有原生延迟消息,可以自己实现一个“时间轮 + Redis ZSet”的延迟队列。Redis ZSet 的 score 存计划执行时间戳,定时任务 zrangeByScore 取出到期任务再投递到业务队列。这个思路在面试里也是加分项。
9. RocketMQ 管理后台与安装排查
很多人在本地部署 RocketMQ 时遇到“mq 安装后管理后台无法进入”的问题。这里给一套通用排查流程。
常见部署方式是先启动 NameServer,再启动 Broker,最后启动管理后台 Console。以 Docker 示例启动的常用模板如下,具体镜像和参数以官方文档为准:
# 启动 NameServer docker run -d --name rmqnamesrv -p 9876:9876 apache/rocketmq:latest sh mqnamesrv # 启动 Broker docker run -d --name rmqbroker \ -p 10911:10911 -p 10909:10909 \ -e NAMESRV_ADDR=rmqnamesrv:9876 \ apache/rocketmq:latest sh mqbroker \ -c /home/rocketmq/conf/broker.conf管理后台无法进入,优先级最高的排查点是这四类:
| 问题现象 | 可能原因 | 排查方式 | 解决方案 |
|---|---|---|---|
| 管理后台页面打不开 | Console 服务未启动或端口未开放 | 检查 Console 容器状态、监听端口 | 确认 Console 启动,映射正确端口 |
| 登录后账号不可用 | 默认账号权限配置问题 | 查看 Console 日志 | 检查 rocketmq-console.properties 或改用明文登录配置 |
| Broker 显示未上线 | Namesrv 地址配置错误 | 进入 Console 查看 Broker 注册状态 | 检查 -e NAMESRV_ADDR 或 broker.conf 中的 namesrvAddr |
| 页面打开但数据为空 | Broker 访问地址使用了容器内 IP | 检查 brokerIP1 配置 | 配置宿主机 IP,如 -e brokerIP1=宿主机IP |
从面试角度讲,管理后台进不去本身不是重点,但它能反映你有没有真正部署过 MQ。面试官喜欢问“你们线上 Broker 挂了怎么处理,怎么通过后台确认积压情况”。你要能说出:通过管理后台看消费者分组、消费位点、消费积压数量,通过 RocketMQ Dashboard 监控告警。
10. 面试高频追问与答题框架
除了基础概念,MQ 事务消息还会被连环追问。这里整理一套高频问题及答案框架。
| 面试官问题 | 答题思路 |
|---|---|
| 事务消息是怎么解决分布式事务的 | 半消息不可见,本地事务执行后提交或回滚,Broker 回查兜底,核心是最终一致性 |
| 为什么不用 2PC | 2PC 强一致但性能差、协调者单点、阻塞资源,MQ 事务消息适合高并发异步场景 |
| 本地事务执行成功,但 commit 失败怎么办 | 依赖 Broker 回查机制,生产者的 checkLocalTransaction 根据业务数据判断 |
| 消费者收到消息后处理失败怎么办 | 重试队列、死信队列、延迟消息补偿、人工处理 |
| 消息重复消费怎么保证幂等 | 数据库唯一键、Redis 判重、业务唯一键 |
| 事务消息会丢消息吗 | 大多数情况下不会,但需要开启同步刷盘、主从复制、Producer 重试等配置 |
| 事务消息和本地消息表区别 | 事务消息由 MQ 实现原子性,本地消息表由数据库和定时任务实现,前者代价更小 |
| 延迟消息能精确到秒吗 | RocketMQ 默认只能指定延迟等级,不能任意秒级;需要精确控制时用定时消息或自研时间轮 |
答题时不要只背流程,要结合订单库存场景,把半消息、本地事务、回查、重试、幂等串成一条线。
如果面试官问“如果回查也失败了怎么办”,你可以回答:最终一致性方案是需要人工兜底的。回查失败的消息会进入事务消息状态异常队列。生产上要加监控,发现事务消息处于 UNKNOW 状态超过阈值就告警,由后台任务或人工处理。这比硬编一个复杂的自动解决方案更实际。
11. 常见问题与排查方法
下面是分布式事务消息生产环境常见的坑排查清单。
| 问题现象 | 可能原因 | 排查方式 | 解决方案 |
|---|---|---|---|
| 消息一直处于半状态,消费者看不到 | 本地事务没返回 commit/rollback,或回查没触发 | 看生产者日志、Broker 事务消息队列 | 检查 TransactionListener 实现,确保回查方法不抛异常 |
| 消费者重复执行业务 | 消费失败后 MQ 自动重投 | 查消费日志、消费位点 | 增加幂等处理,使用业务唯一键 |
| 库存扣了两次 | 订单场景事务消息设计错误,没有幂等 | 查订单与库存流水 | 消费端加状态机,重复消息直接忽略 |
| 延迟消息不按时触发 | 延迟等级设置错,或使用不支持延迟的 MQ | 检查消息的延迟等级和到达时间 | 使用定时消息或自研延迟队列 |
| Broker 重启后消息消失 | 默认同步刷盘/主从配置未开 | 查 broker.conf 中的 flushDiskType、brokerRole | 设置 SYNC_FLUSH,配置主从复制 |
| MQ 管理后台无法进入 | 端口未开放或服务未注册成功 | 检查容器日志、端口映射 | 按前文排查表处理 |
| 消费积压严重 | 消费者性能不足或消费者数量不够 | 看消费位点、消息积压数 | 扩容消费者实例,开启批量消费,优化消费逻辑 |
每条都要能结合自己场景说一句,不要背表。
12. 最佳实践与使用建议
做完上面的设计,再补充几条工程化建议,能明显提升方案完整度。
第一,事务消息的本地事务范围要尽量小。不要在 executeLocalTransaction 里做远程调用、发送短信、调用外部接口。外部调用的失败不应该决定订单事务的提交结果,否则会导致回查链路不稳定。
第二,消费端一定要记录消费流水。每个业务消息都写一条流水,包含 msgId、bizKey、消费时间、处理结果。出现问题时,流水是排查的第一手证据。
第三,重试次数要有限。无限重试会带来消息堆积和重复消费,建议最大重试次数为 3 到 5 次,超过后进入死信队列或人工补偿表。
第四,需要对消息链路做可观测性建设。给消息加 traceId,关联订单号、库存扣减号。日志里打出消息发送时间、消费时间、处理耗时。没有追踪能力的最终一致性方案,出问题只能靠猜。
第五,不要在关键资金链路里只依赖 MQ 最终一致性。如果业务对一致性要求极高,比如资金扣减,优先考虑数据库本地事务和分布式事务框架的结合,MQ 只作为补充异步通知。
13. 总结与下一步
MQ 事务消息和分布式事务的核心,不在中间件本身,而在一致性设计和兜底策略。建议先用 RocketMQ 事务消息把订单与库存场景跑通,再手动模拟 Broker 回查、消费者重复消费、消息发送失败三种异常,观察系统是否能保持最终一致。最容易踩的坑是:本地事务还没提交就返回 COMMIT_MESSAGE,或者消费端没有做业务幂等。先把这个链路理顺,再去扩展延迟消息、批量消费、死信队列、监控告警这些工程能力,面试和实际项目都会轻松很多。