ARTICLE DETAIL

资讯详情

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

MQ事务消息实战:RocketMQ半消息、回查机制与最终一致性

MQ事务消息实战:RocketMQ半消息、回查机制与最终一致性 面试问 MQ 事务消息最怕不是背不出概念而是只能回答一句“RocketMQ 支持事务消息”。真正拉开差距的是你能不能把“半消息、本地事务、回查机制、最终一致性、幂等消费”这一条链路讲清楚能不能现场写出一个订单与库存的示例能不能说出消息丢失、重复消费、延迟消息这些衍生问题怎么处理。这篇文章就把这些点一次性讲透。文中会先给出 MQ 方案在分布式事务里的定位和选型速览然后拆解事务消息的原理再用 RocketMQ 事务消息、本地消息表、幂等消费三个代码示例演示完整设计最后补上延迟消息补偿、管理后台排查、面试高频追问和使用建议。适合准备 MQ 相关岗位面试的开发者也适合正在做订单、库存、支付、积分这类最终一致性方案的工程人员。1. 核心能力速览先说结论MQ 不是用来实现强一致的它是用来实现分布式系统最终一致性的。事务消息是 MQ 在“基础可靠传输”之上提供的一种保证核心目标是“本地事务和消息发送要么都成功要么都失败”。不同 MQ 产品对事务消息的支持程度不一样这里做一张快速对比表对比项RocketMQKafkaRabbitMQPulsar原生事务消息支持支持事务 API用途偏向“精确一次消费”不支持原生事务消息用本地消息表方案支持事务延迟消息支持延迟等级/定时消息不原生支持需自研或用时间轮通过延迟消息插件支持支持定时消息实现复杂度中等消息已封装较高事务 API 理解成本高低但事务场景要改架构中等典型场景订单、交易、支付、库存等电商链路大数据链路、日志、埋点、流处理内部系统异步通知、事件订阅云原生、跨地域复制客户端生态Java 为主多语言完善多语言完善多语言完善面试时说 MQ 事务消息一般默认指 RocketMQ 的事务消息。如果是 Kafka 或 RabbitMQ需要单独说明替代方案比如 Kafka 配合幂等 Producer 和事务 APIRabbitMQ 配合 Publisher Confirm 加本地消息表。2. 前置知识分布式事务的几种典型方案在聊 MQ 事务消息之前必须先铺垫分布式事务的整体框架。否则面试官一追问“为什么不用 2PC”就会卡住。常见分布式事务方案有五类2PC两阶段提交强一致但协调者单点、阻塞时间长性能差。TCCTry-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; MessageString 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 里可以直接拿到。这就是你传订单实体进去的原因。接着实现 TransactionListenerComponent 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) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4;业务侧事务长这样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() { ListMqMessageRecord 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() { ListInventoryAdjustRecord 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_ADDRrmqnamesrv: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 回查兜底核心是最终一致性为什么不用 2PC2PC 强一致但性能差、协调者单点、阻塞资源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或者消费端没有做业务幂等。先把这个链路理顺再去扩展延迟消息、批量消费、死信队列、监控告警这些工程能力面试和实际项目都会轻松很多。
返回列表