RocketMQ 事务消息:Half Message 怎么做到”先发消息、后确认”?
一句话结论(30s)
RocketMQ 事务消息用 Half Message 实现「先发消息、后确认」的二阶段模式,换取最终一致性、替代实现复杂性能差的强一致 2PC。因为 Half 消息写入内部 Topic RMQ_SYS_TRANS_HALF_TOPIC 对 Consumer 不可见,本地事务 Commit 才转写真实 Topic、Rollback 则标记删除,再用回查机制兜底 Commit/Rollback 的网络丢包。
先停一下:为什么事务消息非得「先发消息、后确认」? 因为「发消息」和「本地事务」本质是两次独立操作,没法用同一个本地事务(同一个数据库连接、同一把锁)原子地包起来。如果先发消息再执行本地事务,本地事务失败时消息已投递、消费者已扣款,无法回滚;如果先执行本地事务再发消息,消息发送失败时事务已提交、数据已改,却没人知道要补偿。所以需要「先占坑、后确认」的两阶段方案——half 消息就是那个「占坑但不生效」的第一阶段。
核心原理(2min)
Producer 发 Half Message(Topic 替换为 RMQ_SYS_TRANS_HALF_TOPIC、真实 Topic 存入消息属性)→ 执行本地事务 → 成功 Commit / 失败 Rollback。关键机制:Broker 的 TransactionalMessageCheckService 每 60s 扫描超过 6s 未确认的 Half 消息,回查 Producer.checkLocalTransaction(),返回 COMMIT 转写真实 Topic、ROLLBACK 写入 RMQ_SYS_TRANS_OP_HALF_TOPIC 标记删除、UNKNOWN 下轮再查,最多 15 次后自动 Rollback 防永久悬空;消费端仍需按业务唯一键幂等去重。
底层深入(5-10min)
分布式事务的经典场景
下订单 + 扣库存需要跨服务原子操作。强一致性 2PC 实现复杂且性能差。RocketMQ 的事务消息提供”最终一致性”方案:
1. 发 Half Message(对 Consumer 不可见)
2. 执行本地事务(如扣库存)
3. 本地事务成功 → Commit Half Message(Consumer 可见)
本地事务失败 → Rollback Half Message(消息不投递)
为什么不能直接用强一致 2PC? 因为 2PC 需要协调者在 prepare/commit 两阶段之间锁住所有参与方的资源,锁时间长、协调者单点故障时整个事务悬挂,吞吐被严重拖垮。RocketMQ 用 half 消息把「资源锁定」换成了「消息暂存」,参与者之间不必互相阻塞等待,用最终一致性换性能。
Half Message 存在哪里?
Half Message 不写入目标 Topic。 它写入系统内置 Topic RMQ_SYS_TRANS_HALF_TOPIC(TransactionalMessageUtil.buildHalfTopic() 的返回值就是它):
// TransactionalMessageUtil.java —— buildHalfTopic / buildOpTopic 指向系统内置主题
public static String buildHalfTopic() {
return TopicValidator.RMQ_SYS_TRANS_HALF_TOPIC;
}
public static String buildOpTopic() {
return TopicValidator.RMQ_SYS_TRANS_OP_HALF_TOPIC;
}
Broker 侧真正「换 Topic」的逻辑在 TransactionalMessageBridge.parseHalfMessageInner,先发 Half 时把真实 Topic 塞进属性、再把 Topic 换成 Half Topic:
public PutMessageResult putHalfMessage(MessageExtBrokerInner messageInner) {
return store.putMessage(parseHalfMessageInner(messageInner));
}
private MessageExtBrokerInner parseHalfMessageInner(MessageExtBrokerInner msgInner) {
String uniqId = msgInner.getUserProperty(MessageConst.PROPERTY_UNIQ_CLIENT_MESSAGE_ID_KEYIDX);
if (uniqId != null && !uniqId.isEmpty()) {
MessageAccessor.putProperty(msgInner, TransactionalMessageUtil.TRANSACTION_ID, uniqId);
}
MessageAccessor.putProperty(msgInner, MessageConst.PROPERTY_REAL_TOPIC, msgInner.getTopic());
MessageAccessor.putProperty(msgInner, MessageConst.PROPERTY_REAL_QUEUE_ID,
String.valueOf(msgInner.getQueueId()));
msgInner.setSysFlag(
MessageSysFlag.resetTransactionValue(msgInner.getSysFlag(), MessageSysFlag.TRANSACTION_NOT_TYPE));
if (null != store.getMessageStoreConfig() && store.getMessageStoreConfig().isTransRocksDBEnable() && !store.getMessageStoreConfig().isTransWriteOriginTransHalfEnable()) {
msgInner.setTopic(TransactionalMessageUtil.buildHalfTopicForRocksDB());
} else {
msgInner.setTopic(TransactionalMessageUtil.buildHalfTopic());
}
msgInner.setQueueId(0);
msgInner.setPropertiesString(MessageDecoder.messageProperties2String(msgInner.getProperties()));
return msgInner;
}
Producer 端在发事务消息前,先给消息打上 PROPERTY_TRANSACTION_PREPARED = "true" 标记,再走普通发送流程(DefaultMQProducerImpl.sendMessageInTransaction):
MessageAccessor.putProperty(msg, MessageConst.PROPERTY_TRANSACTION_PREPARED, "true");
MessageAccessor.putProperty(msg, MessageConst.PROPERTY_PRODUCER_GROUP, this.defaultMQProducer.getProducerGroup());
try {
sendResult = this.send(msg);
} catch (Exception e) {
throw new MQClientException("send message Exception", e);
}
分析:parseHalfMessageInner 先把原 Topic 与 queueId 存进 PROPERTY_REAL_TOPIC、PROPERTY_REAL_QUEUE_ID 两个属性,再把 Topic 替换成 Half Topic、queueId 强制置 0。这就是 Half 消息「对 Consumer 不可见」的根源——消费者订阅的是真实 Topic,永远不会去拉 RMQ_SYS_TRANS_HALF_TOPIC。Broker 收到带 PROPERTY_TRANSACTION_PREPARED=true 的请求后走 prepareMessage/asyncPrepareMessage 分支(而非普通 putMessage),从而进入事务状态机。
💭 思考:为什么 Half 消息要「换 Topic + queueId 强制置 0」这一套,而不是简单加个标记位?—— 因为「对消费者不可见」必须靠物理隔离才可靠:消费者只订阅真实 Topic,消息一旦落到
RMQ_SYS_TRANS_HALF_TOPIC就永远不会被拉到,比「打标记、消费时再判断」更彻底、更不可能漏;queueId 置 0 则是把 half 消息统一收敛到同一队列,方便回查线程按序扫描。
Consumer 拉取不到 RMQ_SYS_TRANS_HALF_TOPIC 的消息 —— 因为它们还没有被 Commit。
回查机制:当 Commit/Rollback 丢失
网络超时导致 Producer 的 Commit/Rollback 没有到达 Broker。Broker 定时回查:
1. TransactionalMessageCheckService 定时扫描(默认 60s)
2. 找到 RMQ_SYS_TRANS_HALF_TOPIC 中超过 transactionTimeOut(默认6s) 的消息
3. → 向原 Producer 发送 CHECK 请求
4. → Producer.checkLocalTransaction() 返回:
- COMMIT: 将消息转写到目标 Topic(Consumer 可见)
- ROLLBACK: 消息写入 RMQ_SYS_TRANS_OP_HALF_TOPIC(标记删除)
- UNKNOWN: 回查失败,等待下轮继续回查
回查扫描的核心在 TransactionalMessageServiceImpl.check(),签名直接暴露了两个关键参数:
public void check(long transactionTimeout, int transactionCheckMax,
AbstractTransactionalMessageCheckListener listener) {
...
}
transactionTimeout(默认 6s)决定一条 Half 消息「免检期」结束、进入可回查状态的等待时间,transactionCheckMax(默认 15)决定最多回查多少次。超限判断在 needDiscard:
private boolean needDiscard(MessageExt msgExt, int transactionCheckMax) {
String checkTimes = msgExt.getProperty(MessageConst.PROPERTY_TRANSACTION_CHECK_TIMES);
int checkTime = 1;
if (null != checkTimes) {
checkTime = getInt(checkTimes);
if (checkTime >= transactionCheckMax) {
return true;
} else {
checkTime++;
}
}
msgExt.putUserProperty(MessageConst.PROPERTY_TRANSACTION_CHECK_TIMES, String.valueOf(checkTime));
return false;
}
分析:每轮回查前 Broker 把 PROPERTY_TRANSACTION_CHECK_TIMES 属性 +1 写回消息;当累计次数 checkTime >= transactionCheckMax 时 needDiscard 返回 true,消息被 listener.resolveDiscardMsg() 处理——真实行为是移到系统主题 TRANS_CHECK_MAXTIME_TOPIC 丢弃,而非简单 Rollback。两种归宿本质相同:放弃这条永久悬空的 Half 消息,避免回查线程被它拖死。
最多回查 15 次(transactionCheckMax),超过后消息被丢弃到 TRANS_CHECK_MAXTIME_TOPIC——设计意图是防止永久悬空。
为什么要引入回查? 因为 Commit/Rollback 的 endTransaction 是 oneway 调用,Producer 发出后不等响应,一旦网络丢包或 Producer 进程崩溃,这条 half 消息就永远停在「已发送未确认」状态,消费者永远等不到。回查就是给这种「悬空」状态上的保险——让 Broker 主动去问 Producer「你本地事务到底成了没」。而 UNKNOWN 允许「下轮再查」,是因为本地事务可能还在进行中,不能一查到 UNKNOWN 就武断 rollback。
Commit / Rollback:消息怎么「落地」到真实 Topic
本地事务执行完,Producer 把结果通过 endTransaction 一次性通知 Broker(oneway,不等待响应),只把状态编码进请求头:
switch (localTransactionState) {
case COMMIT_MESSAGE:
requestHeader.setCommitOrRollback(MessageSysFlag.TRANSACTION_COMMIT_TYPE);
break;
case ROLLBACK_MESSAGE:
requestHeader.setCommitOrRollback(MessageSysFlag.TRANSACTION_ROLLBACK_TYPE);
break;
case UNKNOW:
requestHeader.setCommitOrRollback(MessageSysFlag.TRANSACTION_NOT_TYPE);
break;
default:
break;
}
Broker 收到后由 EndTransactionProcessor.processRequest 处理。Commit 的关键一步是 endMessageTransaction——把之前塞进属性的真实 Topic 取回来,重新组装成一条普通消息:
private MessageExtBrokerInner endMessageTransaction(MessageExt msgExt) {
MessageExtBrokerInner msgInner = new MessageExtBrokerInner();
msgInner.setTopic(msgExt.getUserProperty(MessageConst.PROPERTY_REAL_TOPIC));
msgInner.setQueueId(Integer.parseInt(msgExt.getUserProperty(MessageConst.PROPERTY_REAL_QUEUE_ID)));
msgInner.setBody(msgExt.getBody());
msgInner.setFlag(msgExt.getFlag());
msgInner.setBornTimestamp(msgExt.getBornTimestamp());
msgInner.setBornHost(msgExt.getBornHost());
msgInner.setStoreHost(msgExt.getStoreHost());
msgInner.setReconsumeTimes(msgExt.getReconsumeTimes());
msgInner.setWaitStoreMsgOK(false);
msgInner.setTransactionId(msgExt.getUserProperty(MessageConst.PROPERTY_UNIQ_CLIENT_MESSAGE_ID_KEYIDX));
msgInner.setSysFlag(msgExt.getSysFlag());
TopicFilterType topicFilterType =
(msgInner.getSysFlag() & MessageSysFlag.MULTI_TAGS_FLAG) == MessageSysFlag.MULTI_TAGS_FLAG ? TopicFilterType.MULTI_TAG
: TopicFilterType.SINGLE_TAG;
long tagsCodeValue = MessageExtBrokerInner.tagsString2tagsCode(topicFilterType, msgInner.getTags());
msgInner.setTagsCode(tagsCodeValue);
String checkTimes = msgExt.getUserProperty(MessageConst.PROPERTY_TRANSACTION_CHECK_TIMES);
if (StringUtils.isEmpty(checkTimes) && this.brokerController.getMessageStoreConfig().isTransRocksDBEnable() && null != this.brokerController.getMessageStore().getTransMessageRocksDBStore()) {
Integer checkTimesRocksDB = this.brokerController.getMessageStore().getTransMessageRocksDBStore().getCheckTimes(msgInner.getTopic(), msgInner.getTransactionId(), msgExt.getCommitLogOffset());
if (null != checkTimesRocksDB && checkTimesRocksDB >= 0) {
msgExt.putUserProperty(MessageConst.PROPERTY_TRANSACTION_CHECK_TIMES, String.valueOf(checkTimesRocksDB));
}
}
MessageAccessor.setProperties(msgInner, MessageDecoder.string2messageProperties(MessageDecoder.messageProperties2String(msgExt.getProperties())));
MessageAccessor.clearProperty(msgInner, MessageConst.PROPERTY_REAL_TOPIC);
MessageAccessor.clearProperty(msgInner, MessageConst.PROPERTY_REAL_QUEUE_ID);
msgInner.setPropertiesString(MessageDecoder.messageProperties2String(msgInner.getProperties()));
return msgInner;
}
分析:Commit 走「取回真实 Topic → 清除 PROPERTY_REAL_TOPIC/PROPERTY_REAL_QUEUE_ID 事务属性 → sendFinalMessage 写回真实 Topic → deletePrepareMessage 在 OP Topic 打移除标记」;Rollback 则不走 endMessageTransaction,直接 deletePrepareMessage 把 Half 消息标记为已处理,真实 Topic 里不会出现这条消息。这就是「先发后确认」里的确认动作,也是二阶段中第二阶段真正落盘的地方。
💭 思考:为什么 Commit 不是「把 half 消息标记为可见」,而是重新组装一条普通消息写回真实 Topic?—— 因为 half 消息物理上躺在内部 Topic,「标记可见」意味着消费者得去内部 Topic 拉消息,打破正常消费模型;重新组装成普通消息写回真实 Topic,消费者完全无感知,保持了「事务消息和普通消息对消费者一视同仁」的一致性。清除 REAL_TOPIC/REAL_QUEUE_ID 属性则是把消息「洗白」成普通消息,不留事务痕迹。
RMQ_SYS_TRANS_OP_HALF_TOPIC
这个 Topic 记录已处理(Commit 或 Rollback)的 Half Message 的 offset。回查时先查这张表——已处理的直接跳过,避免重复检查和重复投递。
Producer 端的 checkLocalTransaction
TransactionListener listener = new TransactionListener() {
@Override
public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
try {
inventoryService.deduct(arg); // 扣库存
return LocalTransactionState.COMMIT_MESSAGE;
} catch (Exception e) {
return LocalTransactionState.ROLLBACK_MESSAGE;
}
}
@Override
public LocalTransactionState checkLocalTransaction(MessageExt msg) {
// 回查时判断本地事务是否真的成功了
return inventoryService.isDeducted(msg.getTransactionId())
? LocalTransactionState.COMMIT_MESSAGE
: LocalTransactionState.ROLLBACK_MESSAGE;
}
};
上面是用户自己实现的回调。真正调用它的是 DefaultMQProducerImpl.checkTransactionState(Broker 回查请求的客户端入口):
@Override
public void run() {
TransactionCheckListener transactionCheckListener = DefaultMQProducerImpl.this.checkListener();
TransactionListener transactionListener = getCheckListener();
if (transactionCheckListener != null || transactionListener != null) {
LocalTransactionState localTransactionState = LocalTransactionState.UNKNOW;
Throwable exception = null;
try {
if (transactionCheckListener != null) {
localTransactionState = transactionCheckListener.checkLocalTransactionState(message);
} else {
log.debug("TransactionCheckListener is null, used new check API, producerGroup={}", group);
localTransactionState = transactionListener.checkLocalTransaction(message);
}
} catch (Throwable e) {
log.error("Broker call checkTransactionState, but checkLocalTransactionState exception", e);
exception = e;
}
this.processTransactionState(
checkRequestHeader.getTopic(),
localTransactionState,
group,
exception);
} else {
log.warn("CheckTransactionState, pick transactionCheckListener by group[{}] failed", group);
}
}
分析:Broker 的回查请求最终落到 transactionListener.checkLocalTransaction(message) 这个用户实现上。返回值 UNKNOW 会原样透传回 Broker 表示「下轮再查」,回调抛异常则把 remark 一并带回;判断结果再经 processTransactionState 走 endTransactionOneway 通知 Broker 完成 Commit/Rollback。
回查逻辑需要幂等且能独立判断事务状态——不依赖 executeLocalTransaction 的上下文,因为回查可能发生在 Producer 实例重启后。
消费端必须幂等
即使事务消息保证了”本地事务成功才投递消息”,消费端仍可能收到重复消息(网络重试、Producer 回查)。消费端必须根据业务唯一键(订单 ID/事务 ID)做幂等去重。
事务消息已经保证「本地事务成功才投递」,为什么消费端还要幂等? 因为「成功才投递」保证的是「本地事务和消息状态一致」的上限,却挡不住重复投递——回查阶段 Producer 可能把同一事务反复回 COMMIT、Broker 重试、网络重发,都会让同一条业务消息被投递多次。事务消息解决的是「消息与本地事务的一致性」,不是「消息恰好投递一次」,所以消费端必须按业务唯一键去重。
总结
RocketMQ 事务消息用 Half Message 实现”先发后确认”的二阶段模式:Half Message 存在内部 Topic 不可见 → 本地事务执行 → Commit/Rollback 让消息可见/不可见 → 回查机制兜底网络丢包。最终一致性替代强一致性 2PC,是高并发分布式事务的轻量级工程方案。
章末提问
Q1:事务消息为什么要用 half 消息,直接把消息发到真实 Topic 再延迟投递不行吗? 答:不行。因为「发消息」和「本地事务」无法放进同一个本地事务里原子执行。half 消息的本质是「先占坑、后确认」——消息先写到内部 Topic 对消费者不可见,本地事务成功才转写真实 Topic、失败则丢弃,这样消费者要么看到完整成功结果、要么什么都看不到,不会出现「消息投递了但事务回滚」的中间态。延迟投递只是推迟时间,无法解决「事务失败时消息已投递」的脏数据问题。
Q2:回查机制是干什么的?为什么 endTransaction 用 oneway? 答:回查是兜底 Commit/Rollback 网络丢包或 Producer 崩溃的保险。因为 endTransaction 用 oneway 调用(不等响应、追求低延迟),一旦请求丢失,half 消息就永远悬空,所以 Broker 需要定时扫描超过 6s 未确认的 half 消息,主动回查 Producer.checkLocalTransaction() 确认真实状态。oneway 换性能,回查补可靠性,两者是配套设计。
Q3:回查返回 UNKNOWN 会怎样?会不会无限回查? 答:UNKNOWN 表示「下轮再查」,不会无限回查。因为 Broker 每轮回查都会把 PROPERTY_TRANSACTION_CHECK_TIMES 属性 +1,当累计次数达到 transactionCheckMax(默认 15)时,needDiscard 判定超限,消息被丢弃到 TRANS_CHECK_MAXTIME_TOPIC,防止回查线程被永久悬空的消息拖死。
Q4:为什么消费端还必须自己做幂等? 答:因为事务消息保证的是「本地事务和消息状态一致」,不是「消息恰好投递一次」。因为回查可能重复 COMMIT、Broker 重试、网络重发,都会导致同一条消息被投递多次,所以消费端必须按业务唯一键(订单 ID/事务 ID)幂等去重,否则会出现重复扣款。
Q5:half 消息为什么对 Consumer 不可见? 答:因为它根本不写进真实 Topic。因为 parseHalfMessageInner 会把真实 Topic 存进消息属性、再把 Topic 换成系统内部 Topic RMQ_SYS_TRANS_HALF_TOPIC,而消费者只订阅真实 Topic,永远不会去拉这个内部 Topic,所以 half 消息天然对消费者不可见。