先记住这个答案
在 read-process-write 模式中,生产者开启事务,将业务输出和消费位移用 sendOffsetsToTransaction 一起提交;消费者设置 isolation.level=read_committed,只看到已提交输出。此机制只保证 Kafka 内部原子性,外部副作用须用幂等表或 outbox 处理。
- read-process-write 需在同一事务内
- 外部系统副作用须幂等处理
- 隔离级别 read_committed 才看已提交
read-process-write 的原子提交
Kafka exactly-once 的基础是生产者事务。用 transactional.id 初始化后,beginTransaction 到 commitTransaction 之间的多分区写入成为原子单元;abort 后的写入全部不可见。消费者设 isolation.level=read_committed 即可从消费端过滤未提交及已中止的消息。
消费位移可与业务输出绑定:producer.sendOffsetsToTransaction(offsets, groupId) 把已消费 offset 作为内部消息同一事务提交。这样发生重平衡或崩溃时,如果事务已提交则 offset 已持久,未提交则输出被回滚,新实例从原处重读旧数据,避免重复和丢失。
订单解析写入的原子输出
假设输入主题 click_events 有 8 个分区,应用每轮 poll 后做设备解析并写入 parsed_clicks。处理开始时开启事务:先逐条发送解析结果,再对当批记录调用 sendOffsetsToTransaction,最后 commit。若应用在提交前崩溃,事务协调器中止该事务,新消费者因未提交 offset 而重新消费,parsed_clicks 因 read_committed 也看不到任何半成品,从而保证不丢不重。
为避免旧实例残留事务干扰,transactional.id 会生成 epoch 栅栏。同一 transactional.id 下新事务初始化会使旧实例后续事务被拒;生产中启用幂等 enable.idempotence=true、设置 transactional.id 如 'agg-1',并保持 consumer isolation.level=read_committed。这使 read-process-write 在 Kafka 内部成为原子单元,但外部写入不在事务保护内。
事务失效与外部副作用边界
一旦处理期间调用了外部数据库或 HTTP 服务,这些副作用不能被回滚。例如将解析结果写入 Elasticsearch 后 Kafka 事务失败,ES 已留存且无法用 Kafka 撤销。替代方案是给外部写入设计幂等主键,或先用 outbox 在本地事务落库并发布事件,业务表到事件的变更被绑定为原子,再由消费方约束外部操作。
事务还可能因超时或协调器变更失败。默认 transaction.timeout.ms=60000,若某条记录处理超过这个时间,协调器会中止事务而应用可能不知情,导致后续提交被拒绝。通常调大超时并监听 TxnOffsetCommit 异常。另外 Kafka 事务只存在于单个集群内,无法覆盖跨集群的复制链路,因此镜像迁移场景没有 exactly-once。
容易答错的地方
- 幂等生产者等于 exactly-once
- 幂等生产者只解决单分区因重试产生的重复,不解决跨分区原子性和消费者位移一起提交。它不能把消费、处理、写出做成事务,因此单独开启 enable.idempotence 不等于 exactly-once。
- read_committed 保证用户代码只执行一次
- read_committed 只过滤事务可见性。如果用户的 process() 里有写文件或发短信,事务回滚不会撤销这些动作。需要幂等或 outbox 把副作用降级为幂等操作才能在语义上逼近一次。
面试官还会怎么问?
sendOffsetsToTransaction 失败后 offset 会前进吗?
不会,txn 中止则 offset 不提交。进程重启重新拉取,但事务内已发出的消息也被回滚;若部分消息已写入但事务未提交且随后提交成功则可见。需要监控异常并重试整个批次。
只读不写的外部查询算副作用吗?
读取外部状态不产生持久写入,不算必须回滚的副作用。但若读取后基于该状态输出,而外部状态在事务提交前变化,会造成不可重复读。如果要严格一次,需要把外部读取结果也纳入幂等逻辑。
消费者 group.id 与事务如何关联?
sendOffsetsToTransaction 需要传入消费者组 ID 及其已提交的 offsets 映射;事务中写入 offset 后,下次 rebalance 使用该 offset。注意 group.id 必须和消费者的相同,且该消费者不能用 auto.commit。
参考资料
示例用于理解所注明的运行环境与边界;延伸学习可结合原文中的更多案例。