ARTICLE DETAIL

资讯详情

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

Kafka、RocketMQ、RabbitMQ事务消息深度对比:原理、代码与选型避坑

Kafka、RocketMQ、RabbitMQ事务消息深度对比:原理、代码与选型避坑 Kafka、RocketMQ、RabbitMQ这三款消息中间件几乎每个做后端的人都得在自己的技能树里给它们留个位置。但一碰到“事务消息”这个需求很多人会对着三份官方文档陷入困惑Kafka说自己支持事务RocketMQ说自己支持事务消息RabbitMQ也提供了事务API。等你照着各自的思路写完代码跑起来才会发现三者解决问题的范围完全不在同一个层次——Kafka的事务解决的是多分区原子写入RocketMQ的事务消息是为了让业务数据库和下游消费者最终一致而RabbitMQ的channel事务只保证了“消息没丢在生产者到Broker的路上”。这篇文章不为任何一个中间件站台只把三者的原理、代码边界和踩坑经验一次性讲清楚适合正在做MQ选型、准备消息中间件面试或者已经线上踩坑的朋友。1. 先说结论三者口中的“事务消息”根本不是同一个东西1.1 一个每天都在发生的业务难题先还原一个几乎所有业务系统都会遇到的场景。用户下单支付成功后订单服务要做两件事把订单状态改成“已支付”写入数据库给积分服务、库存服务、消息通知服务发一条MQ消息告诉它们“这个订单支付成功了该干嘛干嘛”。问题来了。如果先写数据库再发消息消息发送失败怎么办下游永远不知道这笔订单已支付用户的积分没了库存也没扣财务对账天天报警。如果先发消息再写数据库下游收到消息时数据库可能还没提交下游去查订单状态发现还是“待支付”逻辑直接错乱。更麻烦的是如果发消息和写数据库中间进程崩溃了两边就彻底对不上了。这个问题的本质是“本地事务”和“消息投递”之间没有一致性保证。所谓的事务消息目标就是解决这种不一致要么业务操作和消息投递同时成功要么同时失败或者通过某种补偿机制最终让两边保持一致。1.2 三兄弟对“事务”的不同理解同一个“事务消息”需求三个中间件给出的答案完全不同这是它们出身和设计目标决定的。Kafka出身于大数据领域它的事务是为流处理场景设计的。Kafka Streams要做“读一个topic、处理后写多个topic”这类操作如果写入不同分区的消息一半成功一半失败会造成严重的数据不一致。所以Kafka事务解决的是“Kafka内部多个分区之间的原子写入”它管不到你业务数据库死活。这不是Kafka偷懒而是它压根就不打算把你的MySQL事务纳入管理。RocketMQ出身于电商业务场景。当年的核心诉求就是“业务库事务”和“消息投递”的一致性所以RocketMQ设计了一套半消息机制先把消息发到Broker寄存起来下游看不到等本地事务执行完再通知Broker把消息放出来或者销毁。Broker还会主动回查本地事务的结果防止应用崩溃导致消息一直悬空。RabbitMQ出身于AMQP协议它提供的“事务”是信道级别的txSelect/txCommit/txRollback本质上只保证生产者到Broker之间的消息确认不丢失和你的业务事务没有任何关系。要在RabbitMQ上做业务级事务消息得自己在应用层另起炉灶搭建一套补偿机制。一张表看总体的差异维度KafkaRocketMQRabbitMQ事务核心目标多分区原子写入、流处理精确一次业务事务与消息投递的最终一致信道级消息确认不丢失半消息机制无有无事务回查无有无是否解决DBMQ一致否是否需自行拼装典型使用场景Kafka Streams、多主题原子写电商订单、积分、对账、库存轻量级消息投递2. Kafka的事务做的是流处理里的原子写不是业务库和MQ的对账2.1 诞生背景Kafka Streams的精确一次诉求Kafka在0.11版本引入了事务机制最初的动机就是给Kafka Streams用的。流处理作业有一个经典模式从一个topic读数据经过计算后写入另一个topic甚至同时写多个topic。举个例子一个实时推荐系统从“用户行为日志”主题读取点击流一方面把聚合结果写入“用户画像”主题另一方面把原始行为写入“数据仓库”主题供给离线分析。如果写入过程中Broker挂了或者网络抖动就可能出现一部分消息写到了“用户画像”另一部分消息没写进“数据仓库”。下次作业重启数据就对不齐了。Kafka事务要做到的就是一批消息要么全部写入所有目标分区要么一个都不写入从消费者视角看就是要么全都能读到要么全读不到。这也是Kafka事务最核心的价值它管的范围是“生产者发送到多个分区的这一批消息”和你的订单表、支付表没有任何关系。2.2 底层机制transactional.id、两阶段提交与事务协调器Kafka实现事务的底层机制大体可以拆成几个关键角色。第一个是transactional.id。这是一个你自己定义的事务标识用来在应用重启之后仍然保持事务上下文。设置了transactional.id之后Kafka会隐含开启幂等生产者这样生产者在网络重试时不会产生重复数据。第二个是事务协调器Transaction Coordinator。每个transactional.id都会通过哈希映射到一个内部topic__transaction_state的分区该分区所在的Broker就是这个事务的协调器。协调器负责记录事务的当前状态比如Ongoing、PreparingCommit、Committed、Aborted。第三个是两阶段提交。生产者要commit事务时协调器会先写入一个PrepareCommit标记把事务状态推进到准备提交阶段。然后等所有参与分区的数据都确认写入成功再落一个Commit标记。如果中途任何一步失败协调器就写入Abort标记。消费者侧通过读取这些标记来判断一批消息是否属于已提交事务。这一整套流程里还有个非常重要的细节就是僵尸事务防护。假设你的应用发生网络分区一个旧的生产者实例还活着但新的实例已经接管了同一个transactional.id。如果不做防护两个实例同时写事务数据就乱了。Kafka通过epoch机制解决这个问题——每次新实例注册transactional.id时协调器都会递增事务纪元旧实例再发起请求时会携带旧的epoch协调器直接拒绝。2.3 代码写法和隔离级别用Kafka事务API实现多分区原子写Java代码大致长这样// 设置事务ID同时开启幂等 properties.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, order-pay-transaction-001); properties.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true); KafkaProducerString, String producer new KafkaProducer(properties); // 初始化事务获取或创建一个新的事务纪元 producer.initTransactions(); try { // 开启事务 producer.beginTransaction(); // 向订单主题发送 producer.send(new ProducerRecord(order-topic, order-123, PAID)); // 向积分主题发送 producer.send(new ProducerRecord(points-topic, user-456, 100)); // 两阶段提交协调器会确保两个分区的写入原子生效 producer.commitTransaction(); } catch (Exception e) { // 任何异常都回滚整批事务 producer.abortTransaction(); }消费端也有对应的隔离级别设置。Kafka消费者可以配置isolation.level默认是read_uncommitted事务消息即使没提交也能读到如果设置为read_committed那么事务内的消息只有在协调器落Commit标记后才会暴露给消费者事务失败被Abort的那些消息就永远读不到。2.4 Kafka事务解决不了什么最容易踩的大误区这个误区我在好几个团队的技术方案评审里都遇到过有人想用Kafka事务把“先写MySQL订单表”和“再发Kafka消息”包在同一个事务里以为这样就能保证业务数据和消息的强一致。这个思路完全走偏了。Kafka的两阶段提交只发生在Kafka集群内部MySQL不可能感知到Kafka的PrepareCommit标记Kafka也感知不到MySQL的redo log和undo log。结果就是MySQL事务提交成功Kafka事务commit失败下游依然拿不到消息或者反过来Kafka消息commit成功了MySQL事务回滚了下游收到了一条“假消息”。如果你看到某个教程说Kafka事务能实现本地事务与消息发送的一致那基本是把Kafka事务和RocketMQ事务消息搞混了。Kafka事务的正确用法是把多个topic/分区的写入作为一个整体典型场景是Kafka Streams的精确一次语义或者是“同一个事件需要同时更新两个下游主题副本”这类需求。它没法替你的业务数据库做补偿。另外还有一点Kafka事务是有额外开销的。每开启一个事务生产者要进行FindCoordinator、InitPid、AddPartitionsToTxn等一系列额外的网络交互提交时又有额外的RTT。在高吞吐场景下如果每个业务消息都单独开事务TPS会明显下降。我见过有人把一个批处理任务里的每条消息都单独包事务跑几百万条消息跑了一整晚后来改成“每500条一个事务”速度快了三倍多。3. RocketMQ的半消息与回查只有它把“业务事务补偿”做成了中间件能力3.1 半消息机制先寄存再确认再可见RocketMQ的事务消息在设计之初就是奔着“业务分布式事务”去的这也是国内电商场景倒逼出来的能力。它的核心创新在于半消息官方术语叫Half Message。整个流程是这样的。生产者先不直接发一条普通的业务消息而是发一条半消息到Broker。半消息会被存储到专门的half topic里此刻对消费者完全不可见。然后生产者在自己的本地事务里执行真正的业务操作比如写订单表、扣库存、加积分。本地事务执行完再根据结果向Broker发送Commit或Rollback指令CommitBroker把半消息从half topic转到真正的业务topic消费者开始能看到RollbackBroker直接丢弃半消息。这里的关键点是如果应用在发完半消息之后、执行完本地事务之前崩溃了那这条半消息就会一直悬在Broker上没有Commit也没有Rollback。而RocketMQ之所以被国内大型电商广泛使用就是因为它针对这种情况提供了事务回查能力Broker会定期扫描那些长时间没有确认的半消息然后回调生产者的checkLocalTransaction方法询问“你那个本地事务到底成没成功”。一顿类比就比较直观了半消息就像你到仓库寄存处先登记“我可能有一批货要发”然后你回办公室把账结了再回来通知仓库“这批货正式发出去”。万一你回去路上失忆了仓库会主动打电话问你“账到底结了没有”3.2 executeLocalTransaction与checkLocalTransaction的完整落地姿势RocketMQ事务消息的核心代码有两个回调一个是执行本地事务一个是回查本地事务状态。典型写法如下TransactionMQProducer producer new TransactionMQProducer(); producer.setNamesrvAddr(localhost:9876); producer.setTransactionListener(new TransactionListener() { Override public LocalTransactionState executeLocalTransaction(Message msg, Object arg) { try { // 这里执行真正的本地事务比如写订单表、扣库存 orderService.payOrder(msg.getUserProperty(orderId)); // 本地事务执行成功通知Broker放行半消息 return LocalTransactionState.COMMIT_MESSAGE; } catch (Exception e) { // 本地事务执行失败半消息直接销毁 return LocalTransactionState.ROLLBACK_MESSAGE; } } Override public LocalTransactionState checkLocalTransaction(MessageExt msg) { // 应用崩溃重启后Broker会来回查这里 // 不能查内存不能靠缓存必须查业务库的真实状态 String orderId msg.getUserProperty(orderId); Order order orderService.getOrder(orderId); if (order ! null order.isPaid()) { return LocalTransactionState.COMMIT_MESSAGE; } return LocalTransactionState.UNKNOW; } }); producer.start(); // 第三个参数是自定义业务参数会传给executeLocalTransaction的arg producer.sendMessageInTransaction(new Message(order-topic, PAID), null);落地的时候有几个细节特别容易翻车。第一个executeLocalTransaction里执行DB操作和返回Commit状态不能是同一个事务。如果你在回调里用Spring的Transactional包着函数执行完返回的时候本地数据库事务可能还没真正提交。等你返回Commit给BrokerBroker立刻把半消息变成了消费者可见下游消费者马上来查订单状态结果事务还没提交查到的还是“待支付”一致性直接崩了。正确做法是在步骤A里提交数据库事务确认提交成功后再返回Commit状态。如果数据库事务提交了但返回Commit时网络断了那就等着Broker回查走补偿。第二个checkLocalTransaction必须查数据库不能依赖内存。我之前排查过一个线上事故某团队把事务消息的回查逻辑写成了“查一个内存Map”Map里存了“这条订单是否支付成功”。平时应用不重启看起来一点问题没有。后来一次发版应用重启内存Map清空了Broker回查时所有半消息全部返回UNKNOW回查次数用完直接丢弃。结果几千笔订单在数据库里明明已经是已支付状态下游积分和库存服务就是没收到消息财务对账对了一个星期。内存状态永远不可靠回查必须重新查业务库。第三个消费端依然要做好幂等。RocketMQ事务消息的确减少了“消费到废消息”的概率但不会消除网络重试导致的重复消息。事务消息commit后消费者如果处理成功但ACK丢失Broker会重新投递这就产生了重复。幂等是最后一道防线不管是RocketMQ还是Kafka还是RabbitMQ统统绕不开。3.3 回查参数、幂等要求与和Kafka的本质区别事务回查不是无限次数的。Broker端有几个参数需要重点关注transactionCheckInterval控制回查间隔默认60秒transactionTimeout控制本地事务的超时时间transactionCheckMax控制最大回查次数。UNKNOW状态的消息会被反复回查直到超过最大次数再查不出来就只能人工介入翻数据补消息。这里有一个运维层面的经验回查间隔不要调太短。有人为了让事务消息更快落地把transactionCheckInterval从60秒调成3秒。业务量一大每秒钟Broker要回查大量的半消息每个回查都会打到应用的一个RPC接口纯属回查风暴把下游应用接口打到超时死循环更严重。说到和Kafka的本质区别一句话就能概括Kafka事务的参与者和协调过程都在Kafka集群内部它对外部业务库的最终状态不负责任RocketMQ事务消息的协调范围跨越了Broker和应用通过半消息加上回查机制让Broker能感知业务本地事务的结果。所以同一个“订单支付后通知下游”的需求用RocketMQ事务消息编写完业务库和消息投递都能对得上账。这在Kafka里是办不到的。这也是为什么很多做电商、金融、对账类系统的团队即使已经重度使用了Kafka事务消息这一块还是会单独部署一套RocketMQ。4. RabbitMQ的事务channel级事务和Confirm撑不起业务分布式事务4.1 AMQP事务APItxSelect、txCommit、txRollbackRabbitMQ基于AMQP协议协议层面真真切切提供了事务的能力。在Java客户端里用法是这样的Channel channel connection.createChannel(); try { // 开启事务模式 channel.txSelect(); // 发送消息 channel.basicPublish(exchange-name, routing-key, null, hello.getBytes()); // 提交事务 channel.txCommit(); } catch (Exception e) { // 回滚 channel.txRollback(); }这里的事务保证的其实是“消息已经从生产者发送到RabbitMQ Broker且Broker确认接收了”。等等这不就是confirm机制干的事情吗没错两者目标类似实现方式不同。txSelect是同步阻塞的发一条消息、提交一次事务每个步骤都要等待Broker的确认这种模式在单条消息的确认延迟上还能接受一旦吞吐量上来性能就会断崖式下降。为什么因为每一条消息的发送都伴随着一个事务的开启和提交每次提交都要进行磁盘同步fsync这个开销在消息中间件场景里是灾难级的。我在生产环境做过一个简单的压测对比同样一批消息使用事务模式发送吞吐不到使用普通的自动确认模式的十分之一。所以除非你的系统消息量小到可以忽略否则RabbitMQ事务模式基本就是个展示用的功能不适合业务实际跑。4.2 Publisher Confirm吞吐和可靠性的折中RabbitMQ官方也明确推荐想要保证消息不丢失不要用事务模式用Publisher Confirm机制。Confirm机制是异步的生产者可以批量发送消息Broker确认一批消息之后才回调来通知结果吞吐量远高于事务模式。Channel channel connection.createChannel(); // 开启confirm模式 channel.confirmSelect(); // 批量发送之前调用 long startTag channel.getNextPublishSeqNo(); channel.basicPublish(exchange-name, routing-key, null, msg1.getBytes()); channel.basicPublish(exchange-name, routing-key, null, msg2.getBytes()); // 等待confirm回调或者超时 if (channel.waitForConfirms(5000)) { // 消息都到达了Broker } else { // 部分或全部失败需要补偿重发 }Confirm模式解决了生产者与RabbitMQ之间的消息可靠性问题Broker挂了、网络断了消息没有确认到达应用就可以重发。但它同样不解决“业务数据库事务”和“消息发送”两者之间的一致性问题。你把订单表写入和发送RabbitMQ消息放一起DB事务提交了RabbitMQ发消息时网络超时了下游照样拿不到。4.3 想用RabbitMQ做“事务消息”应用层必须自己造补偿轮子很多团队被业务需求逼着用RabbitMQ做事务消息只因为公司基础设施统一是RabbitMQ不想再引入一套RocketMQ。这种情况下所有补偿逻辑都得自己搭建。业界比较成熟的做法是本地消息表也叫事务消息表或Outbox模式。大致思路是在同一个数据库事务里既写入业务表又插入一条消息表记录两者作为一个本地事务提交。然后由后台定时任务扫描消息表里状态为“未发送”的记录逐条投递到RabbitMQ投递成功再把消息表状态更新为“已发送”。如果投递失败定时任务下次继续扫描重试直到成功为止。这套方案的优点是逻辑直观、不依赖MQ的任何事务能力业务表和消息表天然一致。缺点是应用层要维护那张消息表定时任务要正确设计扫描间隔、批量大小和人工补偿界面。和RocketMQ的中间件级事务消息相比相当于把回查、补偿、状态管理全部搬到了应用层。还有人会在消费端做配合消费者处理完业务后手动ack处理失败时进入死信队列由专门的补偿程序重新投递。但无论如何RabbitMQ本身没有半消息也没有回查那些“先发消息再执行业务失败就回滚消息”的优雅操作在RabbitMQ里完全做不出来。5. 选型判断与避坑经验从面试题到线上故障一条条捋5.1 一张表搞定选型去看各种技术社区的热搜词能看到大量“kafka rabbitmq rocketmq选型对比”“rabbitmq和kafka哪个好用”“rocketmq工作原理”这类问题。选型没有一个绝对标准答案关键看你的业务重心在哪。我给的判断维度很简单判断维度优先选择理由高吞吐、大数据、流处理Kafka分区模型天然适配并行消费事务服务于流处理精确一次电商业务、订单对账、业务事务补偿RocketMQ半消息回查机制把业务事务补偿做成了中间件能力企业内部轻量级系统、已有AMQP生态RabbitMQ部署简单路由灵活管理工具成熟需要强一致而非最终一致以上都不够考虑TCC、Saga等分布式事务框架事务消息只是半路方案5.2 高频坑之一Kafka事务被误用于“DBMQ”一致性这个问题前面已经说了原理但我觉得值得再从面试角度强调一下。面试官问“Kafka事务和RocketMQ事务消息的区别”答案的关键点就是“Kafka事务解决的是Kafka内部多分区写入的原子性不参与数据库事务RocketMQ事务消息通过半消息和回查将业务事务与消息投递的最终一致性问题纳入中间件能力范围。”这样回答后者明显比背文档的人高一层。实际业务中我见过一个团队用Kafka事务做了一个“支付成功通知两个下游主题”的需求设计是对的Kafka事务完美完成了任务。但后来一个新同学不理解硬是在同一个事务里把MySQL更新也塞了进去代码跑了半个月对账不平一查定位到这个事务设计有问题。Kafka事务不是万能事务它就是个“Kafka内部多分区原子写工具”。5.3 高频坑之二RocketMQ回查逻辑写假了消息一直pending这个坑我在3.2里已经详细讲了。再补充一个运维排查链路的经验。线上如果发现消费者一直收不到某类事务消息第一批去看Broker的半消息堆积。RocketMQ控制台或者命令行工具可以看half topic里的消息数量如果半边topic里的消息数量不断堆积说明大量半消息一直没有收到Commit或Rollback。第二步去看生产者应用的日志看checkLocalTransaction是不是被频繁调用以及返回的状态是什么。如果全是UNKNOW再看代码实现多半是查了内存状态或者查了不完整的库。第三步检查回查次数配置如果半消息堆积量大且transactionCheckMax快用完得赶紧人工去数据库核对业务状态手动调用接口或者直接给Broker补发Commit指令避免消息被反复回查后丢弃。这里还要特别注意RocketMQ事务消息的半消息虽然是不可见状态但同样占用Broker的磁盘和内存资源大批半消息悬空会让Broker负载异常升高。5.4 高频坑之三RabbitMQ事务模式的性能雪崩RabbitMQ的事务模式是典型的“看起来很美用起来想哭”。一个真实案例某团队上线了一个千级TPS的订单通知系统开发同学看了RabbitMQ文档里的txSelect觉得简单所有发送都走事务模式上线后RabbitMQ的CPU和磁盘IO直接被打满发送端大量阻塞消息积压到内存。后来排查发现就是每一条消息都经历了一次事务提交和fsync性能从每秒几千条掉到每秒几百条。解决方案是换成了Confirm模式加上批量发送性能恢复了再加上本地消息表做兜底才把业务一致性问题接住。如果你在RabbitMQ上追求“事务消息”记住这条经验事务模式在低吞吐、低频次场景偶可用一旦流量起来它会成为系统里最大的瓶颈。5.5 可以不依赖事务消息的替代方案本地消息表最后分享一个不管选哪个中间件都能用的兜底方案就是本地消息表事务Outbox。在订单表所在的数据库里建一张message_record表字段大致包括消息ID、业务类型、消息体、状态、创建时间、发送时间、重试次数。业务处理时开启本地数据库事务。写入订单状态变更。同时插入一条message_record状态为“未发送”。提交数据库事务。定时任务扫描所有“未发送”且时间超过几秒的记录投递到MQ。投递成功后把状态改为“已发送”删除或保留留痕。这套方案的特殊之处在于消息的可靠性靠的是本地数据库事务而不是MQ的事务能力。就算Kafka或者RocketMQ全部宕机消息也还躺在数据库里系统恢复后定时任务会继续投递。它唯一的代价是应用代码里多了消息表和一套定时任务但稳定性非常高。很多声称“用了RocketMQ事务消息”的系统为了保证最终一致其实也会在应用层加一张类似的表双保险。最后再记一条个人体会。我带新人时经常让他们先忘掉“事务消息”四个字先想清楚两个问题消息能不能重复消费失败能不能重试只要把消费幂等和重试机制做扎实事务消息反而是锦上添花的东西不是雪中送炭的救命稻草。无论你现在打算用Kafka、RocketMQ还是RabbitMQ这句话都适用。
返回列表