系列目录
- MySQL 索引与慢查询:B+ 树如何减少扫描
- 声明式事务之下:InnoDB 的 MVCC 与锁
- 接入层:Nginx 反向代理与 OpenResty 的边界
- Web 容器与 Netty:线程模型之下的 IO 模型
- Redis(上):缓存用法与单线程模型的限制
- Redis(下):超出缓存用途的用法:锁、队列与排行榜
- 数据库访问层:连接池与 MyBatis 的显式 SQL
- Kafka(上):吞吐的来源是顺序 IO
- Kafka(下):生产端、broker 与消费端的可靠性配置
- RPC 框架:像本地调用一样调远程的代价
- 消息语义:按 at-least-once 设计业务代码
- 熔断与限流中间件:把失败当作正常状态管理
- 唯一 ID 中间件 Leaf:号段模式与雪花模式
- RocketMQ 的事务消息与延迟消息(本篇)
本地事务与消息发送的原子性
(Kafka 异步化与削峰)接入 Kafka 做异步化与削峰时,业务数据和待发消息的一致性靠本地事件表解决。做法是把业务写入和一条「待发消息」记录放进同一个本地事务里提交,事务成功则两者都落库,事务回滚则两者都撤销;另有一个扫描任务定时把待发表里的消息投递到 Kafka,投递成功后标记已发送。
本地事件表方案能保证业务写入与待发消息记录同时提交或回滚,但不能让本地事务提交与消息发送作为一个原子操作完成。它分为本地原子写、扫描投递和对账兜底三段。代价包括一张事件表及其清理逻辑、一个扫描任务、由扫描间隔决定的消息延迟、扫描任务故障期间的消息滞留,以及对账补漏。事务回滚后消息不会发出;事务提交后,消息仍可能因扫描任务停止或 MQ 不可用而无法发出,需要监控和对账处理。该方案通过最终一致性处理本地事务与消息发送之间的关系,增加了相应的运维和补偿成本。
2019 年补 RocketMQ 资料时,事务消息提供了另一种实现:broker 参与协调,半消息将本地事务状态与消息对消费方的可见性关联。broker 收到 commit 后才将消息转入真实 topic,消费方此时才能看到消息;收到 rollback 后,消息不会进入真实 topic。这处理的是本地事务提交与消息发送之间的原子性问题。
写作范围:当时交易类业务使用公司 MQ 平台封装的事务能力,没有直接接入开源 RocketMQ 4.3 客户端。下面对事务消息机制的拆解基于 RocketMQ 4.3 官方文档与开源源码;生产侧使用平台封装,内部实现不展开。Kafka(上):吞吐的来源是顺序 IO、Kafka(下):生产端、broker 与消费端的可靠性配置讲了 Kafka 的吞吐来源与可靠性边界,本篇比较 RocketMQ 对业务消息的处理方式。
事务消息:半消息、本地事务、二次确认与回查
RocketMQ 4.3 重新引入事务消息(开源版 4.0–4.2 曾剥离这一功能,4.3.0 于 2018-07 才重新引入,4.3 之前的开源版没有事务消息实现可对照)。机制分四步。
第一步,生产者把消息作为半消息(half message)发给 broker。半消息写入一个特殊的内部 topic RMQ_SYS_TRANS_HALF_TOPIC,不写入业务目标 topic。消费方订阅的是业务 topic,看不到半消息。这一步对应「消息已发出但尚未对消费方可见」的中间态。
第二步,生产者执行本地事务。半消息发送成功后,同一个生产者线程执行本地事务逻辑(扣库存、写订单等)。本地事务的提交或回滚由业务方控制。
第三步,生产者向 broker 提交二次确认。本地事务提交则发 commit,回滚则发 rollback。broker 收到 commit,把半消息从 RMQ_SYS_TRANS_HALF_TOPIC 转投到业务 topic,消费方此时才能消费;收到 rollback,标记该半消息删除,消费方永远看不到。
第四步是回查,用于处理二次确认丢失的情况。生产者执行完本地事务后,可能在发二次确认前崩溃或网络中断,broker 一直收不到 commit/rollback。broker 的回查服务定期扫描半消息 topic 里超时未确认的消息,回调生产者注册的 checkLocalTransaction 接口,查询这条消息对应的本地事务是否已提交。生产者根据本地事务状态返回 commit、rollback 或 unknown。
// 最小示例:事务消息的本地事务执行与回查(脱敏)
TransactionListener listener = new TransactionListener() {
@Override
public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
try {
doBusinessLogic(arg); // 扣库存、写订单等本地事务
return LocalTransactionState.COMMIT_MESSAGE;
} catch (Exception e) {
return LocalTransactionState.ROLLBACK_MESSAGE;
}
}
@Override
public LocalTransactionState checkLocalTransaction(MessageExt msg) {
// broker 回查时调用:业务方根据消息体里的业务键查本地库
String bizKey = msg.getKeys();
return queryTxStatus(bizKey) // 查本地事务是否已提交
? LocalTransactionState.COMMIT_MESSAGE
: LocalTransactionState.ROLLBACK_MESSAGE;
}
};回查需要业务方实现。broker 只能判断半消息是否超时未确认,无法获取本地事务状态;业务方需要查询该状态。回查接口要能根据消息体中的业务键反查本地事务结果,因此本地事务的提交记录必须可查询:要么有事务日志表,要么业务表本身能体现事务是否已提交。如果本地事务提交后没有留下可查询的记录,回查只能返回 unknown,broker 会继续重试回查直到达到次数上限。
回查次数有上限(transactionCheckMax,默认 15 次),超过上限后 broker 默认按 rollback 处理,这条消息被丢弃。如果二次确认丢失且回查持续失败到上限,本地事务已提交的消息仍会丢失,需要业务侧的对账兜底。
事务消息和 消息语义:按 at-least-once 设计业务代码提到的「本地事件表 + 扫描投递 + 对账」都处理本地事务与消息发送的一致性。本地事件表方案由业务侧的扫描任务读取待发表并投递,broker 只做普通投递;事务消息由 broker 维护半消息状态并主动回查。两种方案都要求业务方能查询本地事务状态:本地事件表方案由扫描任务读取待发表,事务消息由回查接口读取本地事务结果。两种方案也都有失败后的处理方式:本地事件表依靠对账,事务消息在回查达到上限后丢弃消息,并由业务侧对账。事务消息在 broker 层根据事务状态控制消息可见性,半消息在中间态对消费方不可见,因此不存在「消息已发但事务还没提交」的可见窗口。本地事件表方案中,扫描间隔决定消费方何时能看到消息。
延迟消息的固定档位和存储开销
生产者发送时设置 delayTimeLevel,消息到点后才对消费方可见。4.x 的延迟级别是固定的 18 档:1s/5s/10s/30s/1m/2m/3m/4m/5m/6m/7m/8m/9m/10m/20m/30m/1h/2h。
延迟档位的数量固定,使 broker 的扫描任务数量和调度频率可以控制。设置了延迟级别的消息不直接写入业务 topic,而是写入按延迟档位组织的内部 topic(SCHEDULE_TOPIC_XXXX,queueId 对应 delayTimeLevel)。broker 的延迟服务(ScheduleMessageService)按档位定时扫描到期消息,到期后将消息重建到真实 topic,再投递给消费方。
固定档位带来两项限制。业务侧只能使用这 18 个延迟值,需要 37 秒延迟时只能选择 30 秒或 1 分钟;任意延迟需要业务侧自行分级或另建方案。每个档位对应一条内部队列,到期投递需要重建消息,档位越细,重建开销越大。固定档位数限制了延迟队列的扫描成本,但也限制了可配置的延迟值。
Kafka 没有内置延迟消息。需要延迟投递时,通常由业务侧的定时任务或外部时间轮实现:消息先写入延迟 topic 或本地表,到点后再投递到目标 topic,成本由业务侧承担。RocketMQ 在 broker 中实现延迟消息,业务侧只需设置一个级别,同时要接受固定档位和消息重建带来的限制。Kafka 将这类功能交给生态或业务侧,broker 以日志吞吐为主;RocketMQ 将常用业务消息功能放入 broker,业务侧的接入工作相应减少。
Kafka 与 RocketMQ 的实现差异
事务消息和延迟消息体现了 RocketMQ 与 Kafka 对 broker 职责的不同划分。本系列第 8、9 篇介绍了 Kafka 的顺序追加分区日志、page cache、sendfile 零拷贝和批量压缩。RocketMQ 在存储模型、消费模型和 broker 内置功能上采用了不同实现。
存储模型。 Kafka 的存储是分区日志:每个 partition 在 broker 上是一组独立的日志段文件,写入追加到该 partition 的活跃段。一个 broker 上有 N 个 partition 就有 N 组日志文件,写入分散到 N 个文件。RocketMQ 的存储是 CommitLog 统一追加:一个 broker 上所有 topic 的所有 queue 的消息顺序写入同一个 CommitLog 文件,再由 ConsumeQueue 按 topic/queue 维度建索引,索引记录每条消息在 CommitLog 里的物理偏移量、大小和 tag 哈希。
CommitLog 统一追加的收益是写入集中在少量文件,磁盘 IO 顺序度更高,在 topic/queue 数量多时写入吞吐更稳定。代价是消费侧按 topic/queue 读取时,需要先查 ConsumeQueue 索引拿到物理偏移量,再回到 CommitLog 取消息,多一次索引跳转;且 CommitLog 混合了所有 topic 的消息,消费某 topic 时数据局部性不如 Kafka 的分区日志。Kafka 的分区日志把同一 partition 的消息物理连续存放,消费侧顺序读对 page cache 友好;topic/partition 数量多时写入分散到多个文件,吞吐会受影响。
CommitLog 统一追加主要优化写入侧,分区日志主要优化读取侧。少量 topic、大流量、按 partition 顺序消费的场景适合 Kafka 的分区日志;大量 topic、每个 topic 流量不大、以业务消息为主的场景适合 RocketMQ 的 CommitLog 统一追加。
消费模型。 Kafka 消费者主动 pull,broker 不维护消费进度(位移存在 __consumer_offsets topic 里,由消费者提交)。RocketMQ 的消费也是 pull,但 broker 端封装成长轮询(push 模式的 API 给业务方,底层仍是 pull),消费进度存在 broker 上。RocketMQ 4.x 还内置了消费重试与死信分级:消费失败的消息自动进入重试队列(%RETRY%消费组名),重试到次数上限进入死信队列(%DLQ%消费组名),重试间隔递增。Kafka 没有内置的重试与死信分级,消费失败的处理要业务侧自己实现(重试 topic 链、死信 topic、对账)。消息语义:按 at-least-once 设计业务代码给出的重试分级加死信加对账是跨 MQ 的通用模式,RocketMQ 把其中重试与死信这部分收进了 broker。
Broker 内置的消息功能。 Kafka 将事务消息、延迟消息和消费重试交给生态或客户端实现,broker 以日志吞吐为核心。RocketMQ 在 broker 中实现事务消息、延迟消息、消费重试与死信,业务侧接入工作减少,broker 复杂度相应增加。功能列表不能单独作为选型依据。日志、数据管道和大吞吐场景适合 Kafka 的存储模型;交易、订单以及需要事务消息和延迟消息的场景适合 RocketMQ 的功能与存储模型。
边界清单
- 消费方仍需处理重复消息。 半消息机制关联本地事务状态和消息对消费方的可见性;消息从 broker 到消费方的投递仍是至少一次,消费侧仍需幂等。消息语义:按 at-least-once 设计业务代码的消费幂等四种实现仍然适用,事务消息不替代它。
- 回查达到上限后消息会丢失。 二次确认丢失时由回查补救,回查次数有上限(
transactionCheckMax默认 15 次),超限后默认 rollback 丢弃。生产者回查接口返回 unknown 或回查本身持续失败时,已提交的本地事务对应的消息可能丢失,需要业务侧对账兜底。二次确认丢失后的未确认状态最多持续到回查上限,达到上限仍未确认时消息会被丢弃。 - 回查要求业务方提供可查询的事务状态。
checkLocalTransaction要能根据消息体反查本地事务是否提交,本地事务必须有可查询的提交记录。业务表或事务日志表要为回查留好查询入口,否则回查只能返回 unknown。 - 延迟消息档位固定。 4.x 只有 18 个固定档位,任意延迟需业务侧自行分级或另建方案。5.x 支持任意延迟,不在本篇口径内。
- 按访问模式选择 Kafka 或 RocketMQ。 Kafka 适合日志、数据管道和大吞吐场景;在少量 topic、大流量下,分区日志存储模型有利于消费侧顺序读取。RocketMQ 适合交易、订单以及需要事务消息和延迟消息的业务场景;在多 topic 下,CommitLog 统一追加使写入吞吐更稳定,broker 内置业务功能减少了接入工作。功能对比表无法说明哪个系统更合适,需要比较业务访问模式与两种存储模型、broker 职责的对应关系。
未深入核对的部分
目前我只了解 NameServer 和 broker 主从切换的配置与使用,没有核对其实现细节。RocketMQ 4.x 用 NameServer 做路由注册与发现(区别于 Kafka 早期依赖 ZK),NameServer 节点之间互不通信,broker 向每个 NameServer 上报路由信息。主从同步有同步复制和异步复制两种模式,broker 角色分 master 与 slave。我没有在源码层核对 NameServer 在节点间不通信时如何保证路由一致性,也没有核对主从切换的具体流程与脑裂风险。ZooKeeper/etcd:小数据强一致的协调服务讲 ZK/etcd 时会补协调服务这一层。
参考资料
- RocketMQ 官方文档(事务消息、延迟消息、存储模型章节,4.x 版本)
- RocketMQ 开源源码(TransactionMQProducer、TransactionListener、ScheduleMessageService、CommitLog 与 ConsumeQueue,4.3+ 版本)
- Kafka(上):吞吐的来源是顺序 IO(Kafka 存储模型与吞吐来源,本篇对照)
- Kafka(下):生产端、broker 与消费端的可靠性配置(Kafka 可靠性配置与位移提交,本篇对照)
- 消息语义:按 at-least-once 设计业务代码(消费幂等与重试死信,事务消息衔接)
- Kafka 异步化与削峰(2016 系列第 6 篇,本地事件表方案,事务消息的对照)
