ARTICLE DETAIL

资讯详情

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

Kafka事务机制详解:两阶段提交、精确一次与避坑指南

Kafka事务机制详解:两阶段提交、精确一次与避坑指南 先聊几句题外话。Kafka 的事务机制很多人把它当成“消息不丢不重”的万能药也有人一听“事务”两个字就想到数据库的 ACID然后按照那个预期去用结果在生产环境踩出一堆坑。我自己在支付流水、库存扣减、订单状态流转这类场景里和 Kafka 事务纠缠了很久今天把这套机制的原理、用法、参数和避坑点一次讲清楚。这篇内容面向需要在实际项目里使用 Kafka 事务的开发者也适合准备面试时想把“精确一次”“两阶段提交”“事务协调器”这些概念真正理解透的人。1. Kafka 事务机制到底解决什么问题先说结论Kafka 的事务机制解决的是“跨多个分区、多个主题的写入原子性”而不是单纯的消息不丢失也不是单纯的消息不重复。它保证的语义是一批消息要么全部写入成功且对下游可见要么全部不可见不存在写到一半的中间状态。1.1 没有事务时我们会遇到哪些问题我在没有事务机制的时期做过一个订单系统流程是订单服务收到请求后往“订单创建”主题写入一条消息再往“库存扣减”主题写入一条消息。正常情况下两个主题的数据是配套的但线上运行一段时间后就会出问题订单主题有条数据库存主题没对应数据或者反过来。原因有几个可能网络抖动导致第一次发送超时但实际已经写入第二次发送时连接断开或者同一个 Producer 在发送两条消息之间进程崩溃重启。这种场景用普通消息发送是无解的。你可能会想两个操作都重试不就行了不行。因为重试本身会造成重复而且两个操作的重试时机不一致依然会出现一边成功一边失败。真正的问题是你没有办法把这两次写入变成一个“要么都成功、要么都失败”的原子操作。Kafka 事务机制干的正是这件事。1.2 事务机制的定位精确一次语义的基石很多人把 Kafka 的事务机制和陈旧概念“exactly-once”绑定在一起但这两者不是一回事。Kafka 事务提供的是“跨多分区原子写入”它和幂等生产者配合才能实现“写入到 Kafka 内部”的精确一次语义。注意这里有个容易混淆的点Kafka 事务只能保证“写入 Kafka”这个动作的精确一次不能保证你的整个业务系统端到端精确一次。比如你从 MySQL 读数据处理后写入 Kafka再用消费者处理写回 MySQL整个过程是否精确一次取决于消费者是不是做了幂等控制、下游有没有配套去重机制。Kafka 事务管不到 MySQL也管不到你的业务代码。这个认知如果不建立后面所有排查都会跑偏。2. 核心机制拆解从 PID 到两阶段提交Kafka 事务的底层设计思路是从老版本的“简单幂等”扩展来的。要理解事务机制先要把几个核心概念搞清楚。2.1 事务协调器与内部主题Kafka 集群里有一个内部主题叫__transaction_state专门记录事务的元数据。每个事务 ID 会通过哈希算法落到这个内部主题的某一个分区上该分区的 Leader 所在的 Broker 会扮演“事务协调器”的角色负责管理对这个事务 ID 对应的所有事务操作。事务协调器是事务机制的大脑。Producer 端发起的initTransactions操作本质是向协调器注册一个事务协调器会生成或校验一个 Producer ID并维护事务的当前状态。后面每次 begin、commit、abort协调器都会记录状态变更并协调各分区的写入。这个设计很容易让人联想到数据库里的“事务管理器”但请注意Kafka 的事务协调器不做数据回滚它只负责记录状态、推进两阶段提交、以及把“事务标记”写入各个数据分区。2.2 两阶段提交的完整过程Kafka 事务采用的是经典的两阶段提交协议但做了很多流式系统层面的优化。整个过程大致是第一阶段是准备阶段。Producer 调用beginTransaction后开始写入数据这些消息会正常写入到各自的分区中但对于使用read_committed隔离级别的消费者来说这些消息是不可见的。这里的关键点是数据已经写进分区的日志文件了只是暂时“不对外出售”。第二阶段是提交阶段。Producer 调用commitTransaction后协调器会先向所有涉及的分区写入一个 prepare commit 标记确认所有分区都准备好提交再写入真正的 commit 标记。当消费者读到 commit 标记时这些事务消息才变为可见。如果是 abort则写入 abort 标记消费者会跳过这些消息。如果用过 RocketMQ 的事务消息你会发现两者的思路完全不同。RocketMQ 是“半消息事务回查”机制消息先被特殊标记等事务提交后才能真正被消费回查机制用于处理极端情况下的未知状态。Kafka 走的是真正的两阶段提交路线代价是事务过程中数据会被写入日志然后通过隔离级别把它们“藏”起来等提交标记出现再释放。2.3 幂等生产者与事务的关系幂等生产者是事务机制的基础设施。开启事务的前提是enable.idempotencetrue内部会自动分配一个 PID并为每个消息附加序列号。Broker 端根据 PID序列号做去重保证同一个 PID 下消息不会重复写入日志。事务机制在幂等生产者的基础上升级了控制粒度。普通幂等生产者只能保证单分区内的消息不重复事务机制将 PID、事务 ID、epoch 结合起来实现了跨分区的原子性和防“僵尸生产者”。所谓僵尸生产者是指某个事务 Producer 发生网络分区或长 GC 后被判定超时但它自己并不知道还在继续发送消息。如果这种情况不处理新生产者会和老生产者产生冲突数据就错乱了。Kafka 通过递增 epoch 来解决新生产者注册事务时 epoch 会加一旧 Producer 拿着已经被淘汰的 epoch 发消息时Broker 会直接拒绝并抛ProducerFencedException。我在实际运维里见过很多次这个报错它表面上看起来很吓人但其实是一种保护机制说明有另一个生产者用了同一个事务 ID 抢占了这个事务的所有权。出现这个异常时正确做法是释放旧的 Producer 实例重新初始化。3. 事务 API 使用与配置要点Kafka 客户端的事务 API 封装得比较清晰但使用顺序、参数配合上到处都是坑。建议大家先把标准流程走顺再考虑扩展。3.1 客户端代码的六步标准流程使用事务性 Producer 的代码流程我总结为六步第一步是构造事务性 Producer。核心配置是transactional.id这个 ID 必须是全局唯一的而且一般要求是一个稳定的标识比如按业务用户维度生成。注意不要在每次请求时都随机生成一个新的 transactional.id那样会反复触发 epoch 重置导致前面的旧事务 Producer 被驱逐。第二步是调用initTransactions()这一步会向事务协调器注册该事务 ID并申请 PID、初始化 epoch。第三步是调用beginTransaction()开启一个新的事务。第四步是正常的发送消息可以往多个分区、多个主题发送所有发送操作都归类到当前事务里。第五步是提交或终止事务调用commitTransaction()或abortTransaction()。最后是在 finally 块里处理异常和关闭资源。有一个非常容易踩的细节beginTransaction()和send()之间不要穿插任何要等待外部系统响应的耗时操作。事务是有超时时间的默认transaction.timeout.ms是 60000 毫秒。如果一件事在事务开启后做了两三分钟才提交协调器大概率已经把事务超时结束了你这边提交时就会收到InvalidTxnStateException或TransactionTimeoutException。3.2 关键参数配置表我整理了一张我常用的配置参考表列了参数名、默认值和建议值大家可以根据自己的场景调整参数名默认值建议值说明transactional.idnull业务前缀稳定标识全局唯一跨重启保持稳定enable.idempotencefalsetrue事务机制强制要求开启acksallall必须全副本确认否则事务无意义transaction.timeout.ms6000060000~120000事务最长允许时间需慎重调整max.block.ms60000可适当调大元数据拉取时被阻塞的时长isolation.level消费端read_uncommittedread_committed事务消息对消费者可见性控制read.isolation.level流处理read_uncommittedread_committedKafka Streams 中的隔离级别配置transaction.timeout.ms这个参数我见过很多人乱调。注意它不是调得越大越好。事务时间越长事务协调器和各分区占用的资源就越久read_committed消费者还会因为等待该事务结束而阻塞下游延迟会飙升。所以这个参数应该根据自己的业务执行时间合理设置而不是随手填一个很大的值。Broker 端还有一个transactions.max.timeout.ms参数默认 900000 毫秒是服务端对客户端事务超时时间的上限如果客户端传的transaction.timeout.ms超过这个值会直接报InvalidTxnTimeoutException。3.3 隔离级别与可见性控制事务消息对消费者不是天然透明的。消费者必须显式设置isolation.levelread_committed才能只读取已提交事务中的消息。如果消费者用的是默认的read_uncommitted那它照样能读到未提交事务中的消息也能读到已回滚事务中被标记删除的数据。这一点经常被忽略结果就是事务机制本身没问题消费者却读到了脏数据排查半天发现是自己的消费端配置漏了。从 Kafka 的存储结构上来理解事务消息写入数据分区时会附带事务标记commit marker 或 abort marker。read_committed消费者在读取某个分区时只会返回那些位于最新已提交事务边界之前的消息。Kafka 内部用 LSOLast Stable Offset这个概念标识这个边界。所有事务消息在 commit marker 出现之前对read_committed消费者来说等同于不存在。我建议所有消费业务数据的消费者默认直接配置isolation.levelread_committed只有明确需要读取实时未提交数据做监控分析的场景才用read_uncommitted。这样能从消费端堵住脏读的入口。4. 实操演示跨分区订单库存原子写入理论讲再多不如上手跑一遍。这里我用一个最常见的业务场景来演示订单创建时同时向“订单主题”和“库存主题”写入数据要求两个主题的数据必须同时成功或同时失败。4.1 场景设计与依赖准备工程依赖只需要 kafka-clients。我这里用 Java 写示例后续的业务逻辑大家根据自己的环境替换即可。先引入依赖dependency groupIdorg.apache.kafka/groupId artifactIdkafka-clients/artifactId version3.4.0/version /dependency准备一个测试用的 Kafka 集群开启事务不需要额外安装组件只要 Broker 配置正常内部主题__transaction_state会自动创建。这里我建议在测试环境使用offsets.topic.replication.factor1、transaction.state.log.replication.factor1来减少资源占用生产环境至少要 3否则内部主题可用性不达标。4.2 完整代码示例与运行效果下面这个 Producer 类是我在项目中用得最多的代码形态大家可以直接参考import org.apache.kafka.clients.producer.KafkaProducer; import org.apache.kafka.clients.producer.ProducerConfig; import org.apache.kafka.clients.producer.ProducerRecord; import org.apache.kafka.common.serialization.StringSerializer; import java.util.Properties; public class OrderTransactionProducer { public static void main(String[] args) { Properties props new Properties(); props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, 192.168.1.10:9092,192.168.1.11:9092); props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true); props.put(ProducerConfig.ACKS_CONFIG, all); props.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, order-tx-connector-001); try (KafkaProducerString, String producer new KafkaProducer(props)) { // 第一步初始化事务申请 PID 并注册事务 ID producer.initTransactions(); // 开启一个新事务 producer.beginTransaction(); try { // 写入订单主题 producer.send(new ProducerRecord(topic-order, order-1001, {\orderId\:\1001\,\amount\:99.9})); // 写入库存主题 producer.send(new ProducerRecord(topic-inventory, sku-5566, {\skuId\:\5566\,\deductQty\:1})); // 如果业务需要可以把消费端的 offset 提交也放入同一个事务中 // producer.sendOffsetsToTransaction(offsetMap, consumerGroupId); // 提交事务协调器向所有分区写入 commit 标记 producer.commitTransaction(); System.out.println(事务提交成功); } catch (Exception e) { // 出现任何异常终止事务 producer.abortTransaction(); System.out.println(事务已终止); throw e; } } catch (Exception e) { e.printStackTrace(); } } }这段代码执行成功后打开任意一个 Kafka UI 工具比如 kafka-ui 或者 AKHQ可以看到两个主题都各自多了一条消息。重点在于如果我把代码改一下比如在第二次 send 之后故意抛一个运行时异常那么这条异常会让流程进入abortTransaction()协调器会向两个主题的分区都写入 abort 标记。此时用read_committed消费者去消费两个主题都看不到任何消息事务的原子性就体现出来了。消费者端的配置如下Properties consumerProps new Properties(); consumerProps.put(bootstrap.servers, 192.168.1.10:9092,192.168.1.11:9092); consumerProps.put(group.id, order-consumer-group); consumerProps.put(key.deserializer, StringDeserializer.class.getName()); consumerProps.put(value.deserializer, StringDeserializer.class.getName()); // 关键只读取已提交事务中的数据 consumerProps.put(isolation.level, read_committed);4.3 如何验证事务真的生效我见过很多人写完代码后直接看消费者有没有收到消息就以为验证过了。真正的验证方法应该是这样的第一步用事务性 Producer 开启一个事务发送消息但既不要 commit 也不要 abort保持事务处于挂起状态。第二步启动一个read_committed消费者你会发现消费不到任何消息等待超时也不会有。第三步将代码恢复为正常 commit消费者立刻就能消费到该事务中的消息。第四步再做一个反向测试让事务回滚消费者仍然消费不到任何消息。如果上面四步都符合预期说明事务机制在你的环境里正常工作了。如果第二步就消费到了未提交的消息那大概率是消费者没配 read_committed或者你用的客户端版本比较老对事务隔离级别支持不完整。另外还可以用 AdminClient 来查看事务状态这在排查问题时很有用try (AdminClient admin AdminClient.create(props)) { ListTransactionsResult result admin.listTransactions(); result.all().get().forEach(t - { System.out.println(事务ID: t.transactionalId()); System.out.println(状态: t.state()); }); }listTransactions会返回事务 ID、生产者 ID、事务状态等信息。配合describeTransactions你还能看到这个事务中涉及了哪些分区、当前处于什么阶段。当你怀疑某个事务卡住了这个命令能直接定位到是哪个事务 ID 在占用资源。5. 常见问题与排查技巧实录我把实际运维中收集到的典型问题整理了一下按出现频率从高到低列一张速查表后面挑几个重点问题细讲。现象可能原因处理方式INVALID_TRANSACTION_TIMEOUTtransaction.timeout.ms超过 Broker 的transactions.max.timeout.ms调小客户端超时时间或调大 Broker 上限ProducerFencedException同一个事务 ID 被新实例抢占旧实例被淘汰终止旧 Producer重新初始化事务消费者读不到已提交事务数据消费端未配置read_committed或生产端未提交成功检查消费端隔离级别配置检查事务状态事务提交超时事务开启时间过长或协调器负载过高缩小事务范围适当调大超时时间TimeoutException或协调器迁移事务协调器所在 Broker 发生故障切换重试提交或检查集群健康状态消息重复事务机制本身就存在重试导致的重复非事务问题需要在下游配合幂等或去重方案事务内消息量过大单事务包含海量消息导致协调器内存压力大控制单事务规模拆分为多个小事务5.1 事务协调器迁移带来的偶发失败Kafka 的事务协调器是根据事务 ID 哈希到内部主题的某分区然后由该分区的 Leader Broker 承担协调工作。如果这个 Broker 宕机、GC 卡顿或发生分区重平衡协调器就会迁移到别的 Broker 上。在此期间Producer 提交事务时可能会出现TransactionCoordinatorFencedException或TimeoutException。这种异常属于一种可恢复异常。业务代码里应该对提交事务做有限次数的重试而不是直接终止整个业务操作。我自己在项目里的做法是在commitTransaction外层加一个最多 3 次的重试循环每次重试前重新执行initTransactions确保新的 epoch 生效。需要注意重试事务提交时可能有部分消息已经被提交了所以业务逻辑里必须保证消息内容的幂等性否则重试会造成重复数据。5.2 消费端多线程场景下如何配合事务热词里有一条是“kafka消费端多线程如何保证消息顺序性”这个问题如果和事务机制一起出现很多人会头大。先说结论Kafka 事务机制和消费端多线程没有直接关系事务只管写入时的原子性不管消费时的顺序性。消费顺序性是一个纯粹的消费端设计问题你要做的是保证同一个业务 key 的消息落到同一个线程里处理。我见过一个项目直接用线程池并行处理消息结果一个订单的“创建”消息被线程 A 处理“支付”消息被线程 B 处理顺序全乱了。解决方案是给线程池提供一个基于 key 的路由策略根据消息 key 的 hash 对线程数取模这样同一个 key 的所有消息会进入同一个队列、同一个线程保持严格的顺序处理。如果你还需要把这个消息的 offset 提交和你的处理结果写到另一个主题的操作做成原子操作那就要用sendOffsetsToTransaction把 offset 提交纳入同一个事务里。这里有个隐藏的坑多线程中每个消费者实例都有独立的 offset 提交需求如果多个线程共享一个事务 Producer又不做严格的并发控制很容易出现事务状态交错导致InvalidTxnStateException。我建议的做法是一个线程对应一个事务 Producer或者对事务 Producer 的调用做线程隔离不要并发调用同一个事务 Producer 的 begin/commit 方法。5.3 消息延迟高与事务的关联热词中“kafka消息延迟高”是很常见的排查项但很少有人会想到事务机制也是延迟的诱因之一。如果某个事务卡在中间状态没有提交也没有回滚那么所有包含该分区数据的read_committed消费者都会被阻塞在该事务的 LSO 之后表现为消费延迟持续上涨。其他事务即使早就提交了但因为日志中前面有一大段未完成事务消费者也只能等。遇到过最夸张的一次是一个开发同学在代码里开启事务后调用了一个第三方的 HTTP 接口把网络超时设成了 30 秒。整个事务链路要等这个接口返回才 commit而且这个接口偶尔还会调用失败不抛异常导致事务永久挂起。结果就是该分区所有消费者全都卡住堆积量一路飙升。排查过程很痛苦最后通过listTransactions才定位到那个“已经开启 40 多分钟”的挂起事务。解决方案是在事务开启之前把所有外部依赖调用都完成事务内部只保留 Kafka 写入操作并且为 commit 操作设置合理的重试策略。如果你确实需要在事务内部做一些外部操作至少把第三方调用的超时时间压到几秒以内并且确保异常路径会走abortTransaction。6. Kafka、RabbitMQ、RocketMQ 的事务能力对比热词里有“kafka、rabbitmq、rocketmq消息队列选型实战对比与避坑指南”既然聊到事务机制这个话题躲不开。很多人面试时喜欢把三个中间件的事务机制总结成“Kafka 有两阶段提交RocketMQ 有半消息RabbitMQ 有事务”这种说法太粗了背后的设计理念和适用场景差别很大。6.1 三类消息队列的事务模型差异RabbitMQ 的事务机制是最“传统”的提供txSelect、txCommit、txRollback这套基于 AMQP 协议的事务能力。它的作用是让一批消息的发送要么全部成功要么全部失败但性能损失非常明显。我在压测里试过开事务比不开事务吞吐量能下降一个数量级所以生产环境中几乎没人用 RabbitMQ 事务来保证原子性更多的做法是用 Publisher Confirm 加业务逻辑兜底。RabbitMQ 事务适合的是单机、低吞吐、强一致的场景在分布式环境下显得很吃力。RocketMQ 的事务消息是另一种思路它引入“半消息”机制和事务回查。生产者先把消息以半消息形式发送到 Broker此时消费者不可见然后业务方执行本地事务执行完后根据结果向 Broker 发送 commit 或 rollback如果网络原因导致这条指令丢失Broker 会定时回查生产者询问事务最终状态。这个模型对我们做分布式事务非常友好特别适合“本地数据库事务消息发送”的经典模式比如订单落库后发消息。Kafka 的事务机制更像是一个“为流处理打造的原子提交协议”它不是为了解决跨系统分布式事务设计的而是为了让 Kafka 内部的多个写入可以被一个消费者以整体方式读取。用 Kafka 事务去协调数据库和消息系统的分布式事务很不顺手因为 Kafka 没有回查机制如果事务提交指令丢失你没办法主动知道这个事务到底提交了没有。6.2 适合选 Kafka 事务的场景与替代方案在消息队列选型这件事上我的建议是这样的。如果你的核心诉求是“本地库表操作和消息发送保持原子”选 RocketMQ 的事务消息会更顺手。如果你追求极致的吞吐又不能容忍消息丢失用 RabbitMQ 的 Confirm 机制配合消费端幂等就够了不一定要上事务。如果你做的是实时计算、流式管道、或者多个 Kafka 主题之间的数据同步那就非常值得用 Kafka 事务。我目前在维护的数据管道就是典型的 Kafka 事务应用场景从原始日志主题读取数据经过清洗和处理把结果写入多个下游主题中间所有写入都由一个事务 Producer 控制。只有这样的设计才能保证多个下游主题能拿到同一批数据而不是有的更新有的没更新。从我个人的维护经验来看Kafka 事务在小规模消息量下没有想象中那么可怕但也不是银弹。它解决了一个清晰定义的问题同时也带来了协调器压力、隔离级别配置、超时处理这些新的复杂度。如果你决定使用它我建议在初期就搭建一套完整的监控手段例如通过 JMX 监控事务指标定期检查__transaction_state内部主题的占用情况并且把事务开启到提交的时间控制在几百毫秒级别不要把它当成一个可以随意包裹长耗时操作的工具。事务范围宁小勿大这是我在多次踩坑后最想强调的一点。
返回列表