先说结论
幂等生产者通过 Producer ID、序列号和 Broker 去重,避免单会话重试造成分区内重复;Kafka 事务可把多分区写入和消费位点提交放进同一事务。消费者使用
read_committed时只读取已提交事务,从而支持 Kafka 内部处理链路的 Exactly-Once。
Exactly-Once 是有范围的语义,不代表发送邮件、调用支付接口或写普通数据库也会自动恰好执行一次。
事务流程
生产者配置稳定且唯一的 transactional.id,初始化事务、发送结果记录,再把来源消费位点加入事务并提交。实例故障后,新的 epoch 会隔离旧生产者,防止僵尸实例继续写入。
Kafka 内事务的完整链路
消费 input-topic offset N
↓
处理并写 output-topic 记录
↓
把 N+1 位点加入同一事务
↓
commitTransaction
↓
read_committed 消费者同时看到输出和新位点
若事务提交失败或生产者崩溃,输出记录对 read_committed 消费者不可见,来源位点也不会推进,重启后可安全重新处理。它将“输出写入”和“输入进度”绑定在 Kafka 内部,而不是把任意外部副作用放进一个魔法全局事务。
幂等与事务的关系
幂等生产者主要解决协议重试导致的同一记录重复写入,范围是一个 producer 会话与分区;事务在此基础上组织多个分区写入、消费位点提交和失败隔离。开启事务不意味着应用主动执行两次业务 send 会自动合并,两次不同的业务操作仍是两条记录。
transactional.id 必须对同一逻辑生产者稳定且在并发实例之间不冲突。新实例使用相同 ID 会提高 epoch 并 fenced 旧实例,避免网络分区时两个“同一身份”同时写入。
消费者隔离级别
默认消费者可能读取未提交事务写入的数据边界;设置 isolation.level=read_committed 后会跳过已中止事务并等待未完成事务的稳定边界。这个选择会影响延迟和可见性,整个需要 Exactly-Once 的下游链路都应保持一致的隔离设置。
外部系统为什么仍需 Outbox
数据库事务提交成功
↓
Outbox 表记录 order.paid 事件
↓
CDC/投递器发布 Kafka
↓
消费者以事件 ID/版本幂等落库
Kafka 事务无法原子提交 MySQL 的订单更新和 Kafka topic 写入。Outbox 把“业务状态”和“待投递事件”放到同一数据库事务,之后通过可重试的异步投递收敛;消费者同样需要唯一约束或状态机处理重复。
性能和运维成本
事务维护协调器状态、超时、marker 写入和额外 RPC,会增加延迟和资源使用。应只用于需要 Kafka-to-Kafka 原子处理的链路,并监控事务超时、abort、fencing、未完成事务和 broker 资源。长事务会阻塞 read_committed 消费进度,事务范围应尽量短。
外部系统一致性
写数据库并发 Kafka 消息通常采用 Outbox、CDC 或业务幂等。消费者落库使用唯一键、版本状态机和事务,把“处理结果”和“去重记录”原子提交。
容易踩坑的地方
事务会增加请求、延迟和状态管理成本,不应用于所有普通发送。消费者若使用默认隔离级别,可能看到未提交或已中止事务相关的数据边界,必须按链路要求配置。
常见问题
追问:幂等生产者能防止应用重复调用 send 吗?
不能。它主要去除协议重试的重复;应用主动发送两次不同记录仍会写入,业务重复需要业务键和下游幂等处理。
追问:Exactly-Once 是否等于消息只处理一次?
不是。处理代码可能重跑,Exactly-Once 描述的是在规定范围内对可见输出的等价效果。跨数据库、邮件、支付等副作用仍需要业务幂等和补偿。
追问:事务超时会发生什么?
协调器会中止长时间未完成事务,生产者后续提交失败,需要终止当前处理并重新初始化或恢复。事务内不能包含长时间阻塞操作。
追问:read_uncommitted 适合什么场景?
不要求事务隔离、追求最低延迟或用于诊断时可使用;但它可能看到之后被中止的记录,业务必须能容忍或自行过滤。