这些年我带着团队做过不少基于 Kafka 的数据管道项目,从最简单的日志投递到订单状态机流转,几乎每一轮都会有人问:Kafka 到底能不能做到精确一次消费?所谓的“原子更新”又是什么?坦白讲,很多人把 Exactly-Once 当成一个“配置开关”,以为改一个参数就能点石成金。实际上它是一整套机制的组合——幂等生产、事务协调、消费侧隔离级别,再加上业务侧的状态幂等,缺一环都白搭。这篇笔记就按我自己的学习路径,把 Kafka 原子更新和精确一次消费(Exactly-Once)从原理到代码再到线上排障,完整过一遍。
1. 先别急着配参数,三种投递语义必须盘明白
1.1 At-Most-Once、At-Least-Once、Exactly-Once 到底在说什么
消息系统最绕不开的三个名词,初看像英语考试,实际上对应的是“丢、重、稳”三种业务感受。
At-Most-Once 是“最多一次”,一条消息可能被处理一次,也可能一次都不处理。Kafka 里如果生产端把 acks 设成 0,发完就认为成功,或者消费端先提交 offset 再执行业务逻辑,一旦崩溃,消息就丢了。语义上叫 at-most-once,现实中往往叫“数据缺失事故”。
At-Least-Once 是“至少一次”,不丢,但可能重复。这是 Kafka 默认消费方式最常见的结果:消费者拉取一批消息,处理完了,还没来得及提交 offset,进程重启,rebalance 之后又把这批消息重新拉了一遍。业务如果没做幂等,就会出现重复扣款、重复发券之类的现场。
Exactly-Once 是“精确一次”,既不多也不少。这个目标在单机单库场景很容易,靠数据库事务就能搞定;到了跨网络、跨节点、跨消息系统的分布式环境,就成了需要一层一层机制堆出来的硬骨头。
1.2 “精确一次”难在分布式
为什么单机数据库能直接靠事务,Kafka 却绕了这么一大圈?核心原因只有一个链条,叫“状态不一致窗口”。
生产者把消息发到 broker,消息落盘成功,但 ack 在半路上丢了,生产者不知道成功,选择重试,于是 broker 收到了两条一模一样的数据。这是“确认丢失导致重复”。
消费者把业务逻辑执行完了,结果把当前 offset 提交到 Kafka 的动作失败或来不及做,进程一挂,重启后重复消费一整批。这是“偏移量提交和业务执行不同步导致重复”。
更麻烦的是,一条业务消息可能要写入多个分区、甚至多个 topic,broker A 写成功了,broker B 写失败了,外部系统就会看到一个“只完成了一半”的更新结果。对这种多个独立节点之间发生变化的过程,如果没有一个“要么全做、要么全不做”的裁定者,任何链路都是不可靠的。
所以 Exactly-Once 要解决的从来不是单个网络包的可靠性,而是“在不确定环境下,让多个参与方就同一批数据的终态达成一致”。
1.3 什么业务必须吃这套
不是所有场景都需要精确一次,强上事务反而拖吞吐。我自己的判断标准是三条:涉及钱、涉及状态、涉及对外对账。
- 支付和订单:重复扣款不可接受;状态机从 PENDING 到 PAID 只能推进一次。
- 库存和券类资源:扣减和发放必须只执行一次,容不得超卖。
- 审计流水和积分系统:多一条流水,月底对账直接爆炸。
反观日志采集、行为埋点这类场景,漏一批顶多统计分析少几个点,重复几条也不影响结论,完全没必要开事务给自己找罪受。
2. Kafka 的“原子更新”到底更新了啥
2.1 幂等生产者:第一层防重复
先说明一点,幂等生产者解决的是“一个生产者会话内,向单个分区发送时不重复、不乱序”。它不解决跨分区的事务问题,更不解决消费端重复。
开启方式非常简单,enable.idempotence=true,在 Kafka 3.0 之后这已经是生产者的默认值了。但它的内部机制值得记一下,面试和排障都用得上。
每个生产者启动时会从 broker 申请一个 PID(Producer ID),然后向每个分区发送消息时,都会带一个从 0 开始单调递增的 sequence number。broker 端会维护一张“PID + 分区 + 序号”去重表,收到一条新消息,先看序号是不是自己期望的下一个。如果重复,直接丢;如果序号跳变,说明中间真有消息丢了,直接报 OutOfOrderSequenceException。
这有点像一个叫号系统:医生叫到 31 号,来了一个病人说自己是 31 号,直接看;又来一个还说 31 号,护士一看叫过号了,拒绝重复接诊。但问题是,如果医生换了一个人,原来的叫号记录全没记住,那 31 号可能又被看一遍。
这就是幂等生产者的边界:PID 一旦变化(进程重启、broker 端状态丢失),去重记忆就断了。所以幂等只能在“单会话”内生效,想要跨会话、跨分区的完整性,必须上事务。
2.2 事务机制:让多分区写入具备原子性
Kafka 的事务,本质上干的就是“原子更新”这件事——把多个分区、多个 topic 的消息写入,包装成一个原子操作。要么所有消息都对外可见,要么全部不可见,不存在中间状态。
实现起来,Kafka 引入了一个新的协作角色,叫 Transaction Coordinator(事务协调器)。每个配置了transactional.id的生产者,都会通过哈希算法被分配给某一个协调器负责,协调器的状态会落到一个内部主题__transaction_state里。
一次完整的事务流程是这样走的:
- 生产者调用
initTransactions(),向协调器注册事务 ID 和新的 PID,获取一个递增的 epoch 值,用来“屏蔽”旧生产者,防止僵尸进程干扰。 - 生产者向目标分区发送业务消息,此时消息只是“写入但未提交”,对 read_committed 消费者不可见。
- 生产者调用
commitTransaction(),协调器在事务日志中记录提交状态,然后向每个涉及的分区写入一个 COMMIT 控制消息。 - 消费者看到 COMMIT 标记,才把事务范围内的消息暴露出来。
反过来,如果任何一步出错,生产者调abortTransaction(),协调器会向各个分区写入 ABORT 控制消息,消费者会把事务内写入的消息全部屏蔽掉。
用餐厅打个比方:一个包间有十道菜,厨房不可能做一道上一道、客人吃一道走一道。后厨全部做完,领班确认“这桌上齐了”(COMMIT),才能一起端出去;如果中途发现某道菜做坏了,整套菜全部撤掉(ABORT),客人绝不可能吃到九道半的菜。
2.3 消费侧的隔离级别:read_committed 是怎么做到“不见不发”的
很多人开了事务生产者,消费端却不改配置,结果发现照样读到未提交消息、甚至读到被回滚的消息。这是因为消费者有两个隔离级别:
read_uncommitted:默认值,能读到事务内已写入但未提交的消息,也能读到已 abort 事务写入的数据。read_committed:只能读到已提交事务的消息,且自动跳过 abort 事务的数据。
但这里有一个很多资料没讲透的细节:read_committed 消费者不是简单做个过滤,它会被一个叫 LSO(Last Stable Offset,稳定末尾偏移)的概念卡住。
Kafka 分区里消息是线性存储的。如果一个长事务迟迟不提交,它的消息就卡在分区中间。read_committed 消费者只能读到 LSO 以内的部分,LSO 之后的已提交消息,哪怕物理上已经存在,也会被挡住,直到挡路的事务提交或回滚,消费才能继续往前。
这个机制保证了原子可见性:绝对不会出现“我看到了 topic A 里的消息,却没看到同一个事务写入 topic B 里的配套消息”这种半成品。
3. 实操:把 Exactly-Once 从配置到代码拉通
3.1 集群侧配置清单
学习阶段别用单节点做事务实验,否则第一个报错就能劝退你。Kafka 事务依赖__transaction_state内部主题,默认副本数是 3,最小 ISR 是 2。单节点集群根本满足不了,会直接报:
Number of alive brokers '1' does not meet the required replication factor '3''本地实验可以临时把这些参数调低,但生产环境我强烈建议保持在默认陪数,这是事务状态高可用的底线。
- broker 配置:
transaction.state.log.replication.factor=3、transaction.state.log.min.isr=2、transactions.max.timeout.ms=900000。 - 生产者配置:
transactional.id(全局唯一)、enable.idempotence=true(事务开启后默认会启用)、transaction.timeout.ms=60000。 - 消费者配置:
isolation.level=read_committed、enable.auto.commit=false。
如果你还在用 Zookeeper 模式,记得把三台 broker 起全;用 KRaft 模式也一样,控制器节点之间也要保证多数派。集群没起全就开事务,后续问题排查会非常痛苦。
3.2 事务性生产者代码与运行结果
Java 客户端是最常用的,直接用 KafkaProducer 原生 API 就能跑通。先看完整代码:
Properties props = new Properties(); props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "node1:9092,node2:9092,node3:9092"); props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer"); props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer"); props.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, "order-tx-producer-01"); props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true); KafkaProducer<String, String> producer = new KafkaProducer<>(props); producer.initTransactions(); try { producer.beginTransaction(); producer.send(new ProducerRecord<>("order-topic", "order-1001", "transaction-data-1")); producer.send(new ProducerRecord<>("order-status-topic", "order-1001", "PAID")); producer.commitTransaction(); } catch (Exception e) { producer.abortTransaction(); }这段代码最核心的地方在于:beginTransaction()把两个不同 topic 的写入圈进同一个事务,commitTransaction()之前如果任何一次 send 失败,都会走 abort,保证两个 topic 各自都不会出现“只见一半”的状态。
运行过程中的经验:
transactional.id一旦设置,必须全局唯一。两个生产者抢同一个 ID,会出现 ProducerFencedException,这是 epoch 机制在主动保护数据,不是 bug。- 初次
initTransactions()会比较慢,因为要和协调器建立会话、申请 PID。不要在构造函数里做重量级事情,程序启动阶段就初始化好。 - send 返回的是 Future,建议用回调捕获异常,别用
get()把每条消息都阻塞一遍,事务本身已经兼顾了可靠性,业务线程没必要再串行化。
3.3 事务性消费者的正确食用方式
消费者侧要做两件事:把隔离级别设为 read_committed,把自动提交关掉,自己控制 offset。
Properties props = new Properties(); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "node1:9092,node2:9092,node3:9092"); props.put(ConsumerConfig.GROUP_ID_CONFIG, "order-consumer-group"); props.put(ConsumerConfig.ISOLATION_LEVEL_CONFIG, "read_committed"); props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false"); props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props); consumer.subscribe(Collections.singletonList("order-topic")); while (true) { ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(200)); for (ConsumerRecord<String, String> record : records) { process(record); } consumer.commitSync(); }但这里我要泼一盆冷水:原生消费者 API 的commitSync()提交 offset,和事务生产者的 commit 事务,不是同一个原子操作。也就是说,就算你开了 read_committed,也绝不可能做到“业务处理 + offset 提交”在同一个事务里原生完成。
Kafka 的端到端精确一次,在原生 API 层面其实是个伪命题,真正能把 offset 放进事务一起提交的,是 Kafka Streams 这类上层框架。普通项目更现实的姿势是:用事务生产者做原子写,消费端自己做幂等;或者走 3.4 里的 Outbox 模式。
3.4 业务原子更新的三种落地套路
如果你要的不是“Kafka 内部消息的原子性”,而是“业务数据库状态更新 + 发消息”的原子性,那么请记住三种模式,由轻到重。
模式一:消费端幂等。在业务表里建唯一键,比如订单号。重复消费时插入冲突,update 不生效,数据不会翻倍。这是最简单也最常见的兜底方案,但对“状态机多次推进”就帮不上忙了。
模式二:事务性消息 + 状态校验。消费到消息后,更新业务库,用乐观锁(比如版本号)挡住重复消息。第一次更新成功,后续重复消息会因为 update 影响行数为 0 而直接跳过。
模式三:Outbox 模式。把业务写库和通知消息塞进同一个数据库本地事务:业务变更和 outbox 表一起提交,再通过 CDC 或后台轮询把 outbox 里的记录发布到 Kafka。下游消费者收到的是“已经经过业务事务确认”的可靠消息,数据库和消息系统的原子性由数据库本地事务来保证。
其中 Outbox 模式,我愿称之为“用数据库事务给 Kafka 补原子性”的最优雅方案。比如你更新订单状态为 PAID,同时往 outbox 表插一条消息,这两个操作在一个数据库事务里提交,要么都成,要么都败。之后消费者把消息发到 Kafka,就算失败也可以重试,outbox 记录不会消失。这个套路在微服务架构里比强行追求跨系统事务要靠谱得多。
4. 线上踩坑与排查速查
4.1 直接报错“Number of alive brokers”是什么鬼
这个报错我见过太多次了,十有八九是单节点/双节点环境想开事务。事务状态日志的每个分区要求 3 个副本,双节点只有 2,达不到 ISR,协调器就无法写入状态。
开发环境解决办法有两个:一是起 3 个 broker,最省事;二是修改 broker 配置,把transaction.state.log.replication.factor和transaction.state.log.min.isr调低,但生产环境绝对不要这么干。
还有一个隐藏坑:如果你的集群是旧的 Zookeeper 模式,升级到新版本后,__transaction_state的副本因子可能需要手动调整。检查命令还是老一套:
kafka-topics.sh --describe --topic __transaction_state --bootstrap-server node1:9092看到Isr齐了、ReplicationFactor达标,再往下查。
4.2 消息延迟高,第一个怀疑对象是事务
开事务之后延迟变高,是正常的,本质是交易的成本。最明显的变化有两个:
第一,生产者在提交事务时要和协调器多几轮网络交互,协调器又要往事务日志分区写状态,整体 RTT 变多。第二,read_committed 消费者会被 LSO 挡住,如果前面长期挂着一个未提交的大事务,后面的已提交消息全被压在分区里。
排查思路从三个方向入手:
- 看有没有“僵尸长事务”。在 broker 日志里搜
InitProducerId、EndTransaction相关的记录,看是否有事务长时间没有 end。也可以用kafka-run-class.sh kafka.tools.TransactionsTool之类工具去观察事务状态。 - 看协调器负载。事务协调器是按
transactional.id哈希分布的,如果一个协调器同时管了太多大事务,状态写入会变慢,进而拖慢所有关联生产者。 - 判断是否被 LSO 卡死。消费者指标里有
records-lag,如果 lag 持续不变且时间戳不再推进,配合 JMX 或分布式追踪看当前 partition 是否存在 pending 事务。
另外一个优化经验:大事小做。把一次事务里的消息量控制在合理范围,不要一个事务塞几百万条。Kafka 事务本身是轻量级协调机制,但太大就会造成“提交时间长、阻塞消费久”的连锁反应。
4.3 提交了事务,消费者为什么就是 poll 不到
有一个非常经典的现象:生产者commitTransaction()已经返回成功了,业务日志也打了,read_committed 消费者却一直 poll 不到数据。
我第一次遇到时也懵了,后来查明白了,原因就是 LSO 卡顿。消费者 poll 到的位置不能越过“第一个未决事务”。假设分区里前半段有一个消费者启动很早之前就存在的事务,因生产者进程被杀一直没有 commit,协调器在超时前也没来得及回滚,那这个分区上 LSO 就停在那里,后面即使堆了几十条已提交消息,read_committed 消费者也只能眼睁睁看着 lag 不下降。
处理办法:找到挂掉的生产者对应的事务,要么等超时回滚,要么手动干预事务状态。在运维层面,保证每个使用事务的实例都有完善的优雅停机逻辑,别让进程被 kill -9 后留下大量孤儿事务。
还有一种情况是消费者把isolation.level配错。默认是read_uncommitted,如果你只开生产者事务、没改消费者配置,会读到 abort 事务里的半成品数据,这也会造成“数据看起来很奇怪”的现场。
4.4 面试必问的几条辨析
这些辨析我整理成了一个速查表,背过之后至少面试不会露怯:
| 对比维度 | 幂等生产者 | 事务机制 |
|---|---|---|
| 解决范围 | 单生产者会话内、单个分区内不重复 | 跨会话、跨分区、跨 topic 的原子写入 |
| 核心机制 | PID + sequence number 去重 | Transaction Coordinator +__transaction_state |
| 失败回滚能力 | 无,只去重 | 有,abort 后 check_committed 消费端自动跳过 |
| 对消费端要求 | 无 | 建议设 read_committed |
| 性能开销 | 较低 | 较高,多协调轮次和状态写入 |
再补一个高频题:acks=all 和 Exactly-Once 的关系。acks=all只保证消息不丢,不保证消息不重。因为生产者发送成功后如果 ack 丢失,还是会重试,broker 收到重复批次。幂等生产者解决的就是这个重复批次问题,而事务在幂等基础上又解决了多分区原子性和生产者会话切换后的僵尸写问题。
最后还有一个大家常问的:Kafka 到底有没有 UI 界面?答案是有的,社区里常见的 Kafka UI(provectus/kafka-ui)、Kafka Manager(CMAK)、Kowl 这类工具都可以看 topic、消费组和消息详情。但说句实话,事务协调这块的深度排障,UI 只能帮你定位现象,最终还是要回到命令行工具和 broker 日志上去。
5. 最后分享一点个人体会
纸上谈兵这么多,落到实际项目里,我的选择标准特别简单:如果是订单、支付、库存这类链路,我会开事务生产者 + read_committed 消费者,并且要求业务方对关键操作做幂等,两头一起堵;如果是日志和埋点,绝不开事务,幂等都不开,直接 acks=1 换吞吐。
学习 Kafka 事务最好的办法,是自己搭一个三节点集群,写一个事务生产者,再手动 kill 掉生产者的进程,观察一个未决事务怎么被协调器回滚,看看 read_committed 消费者在 LSO 前后数据可见性发生了什么变化。把这个实验做一次,你对“原子更新”和“精确一次”的理解会超过看十篇文档。还有个小习惯,写事务代码时逻辑很绕,推荐在一个事务开始和结束的地方打上清晰日志,线上定位问题能省很多力气。就是要记住:精确一次不是一个开关,是一整套系统的共同承诺。