这是《Kafka 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。 “消息会不会丢、会不会重”是 Kafka 工程的灵魂问题。本章把三种交付语义的定义、实现手段与适用边界彻底讲清楚,并给出可以落地的代码。
14.1 三种交付语义
| 语义 | 含义 | 典型实现 | 风险 |
|---|---|---|---|
| At most once | 至多一次 | 先提交位移,后处理消息 | 可能丢,不会重 |
| At least once | 至少一次 | 先处理消息,后提交位移 | 可能重,不会丢 |
| Exactly once | 精确一次 | 幂等 + 事务 + 消费端幂等存储 | 实现复杂,范围有限制 |
注意一个常见误解:Exactly once 不是魔法,而是若干机制组合出的端到端属性,而且它的“端”有明确边界。Kafka 事务保证的是“Kafka 到 Kafka”的处理链路精确一次;一旦涉及外部系统(MySQL、Redis),仍需要业务侧幂等或事务配合。
14.2 至少一次:默认的工程底线
标准做法(第 6 章):
poll -> 处理 -> commit offset
处理完成但未来得及提交时宕机,重启后会重复消费,所以是“至少一次”。重复交给下游幂等解决:
- 数据库唯一键
INSERT ... ON DUPLICATE KEY UPDATE; - Redis
SETNX key requestId; - 状态机检查(订单已是 PAID 则忽略再次支付事件);
- 业务版本号 / 乐观锁。
绝大多数业务系统应该停在这里:至少一次 + 业务幂等。 简单、可控、可观测。
14.3 幂等生产者:解决发送端重复
问题:为什么发送会重复
Producer 发消息给 Broker,Broker 写入成功,但响应在网络中丢了。Producer 超时重试,Broker 收到两条相同消息。网络层面无法区分“没写进去”和“写进去了但响应丢了”。
原理:PID + Sequence Number
开启幂等(enable.idempotence=true,3.0 起默认开启)后:
- Producer 启动时从 Broker 获取一个 PID(producer id);
- 对每个
(PID, partition)维护单调递增的 sequence number; - 消息(batch)带着
PID + epoch + sequence写入 Broker; - Broker 为每个分区缓存最近 5 个 batch 的序号:
- 序号正好是预期值:正常写入;
- 序号小于预期(重复):丢弃并返回成功;
- 序号大于预期(乱序缺口):返回
OutOfOrderSequenceException。
Producer 重试: seq=7 -> Broker 已有 seq=7 -> 丢弃,返回成功
客户端视角:发送成功,且日志中只有一条
约束与注意:
- 幂等范围是当前 Producer 会话 + 单分区,进程重启后 PID 改变,不能防止跨重启的重复;
max.in.flight.requests.per.connection <= 5且幂等开启时,Broker 能重排序号,仍保持分区内有序;- 若关闭幂等又想要顺序,只能把 in-flight 设为 1,吞吐大幅下降。
14.4 事务:跨分区原子写
幂等解决“单分区单会话不重”,事务解决“多条消息原子写入多个分区”。
完整 API
Properties props = new Properties();
props.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, "order-tx-producer-1");
props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true);
// 其余配置同普通生产者
KafkaProducer<String, String> producer = new KafkaProducer<>(props);
producer.initTransactions(); // 1. 注册到事务协调器,获取 epoch(会 fence 旧实例)
try {
producer.beginTransaction(); // 2. 开启事务
producer.send(new ProducerRecord<>("order-events", "o1", "created"));
producer.send(new ProducerRecord<>("audit-events", "o1", "audit"));
producer.commitTransaction(); // 3. 提交:两条要么都可见,要么都不可见
} catch (ProducerFencedException | OutOfOrderSequenceException | AuthorizationException e) {
producer.close(); // 不可恢复,直接关闭
} catch (KafkaException e) {
producer.abortTransaction(); // 可恢复错误则回滚
}
要点:
transactional.id必须稳定且唯一(如“服务名-分区号”)。重启后新实例用同一 id 初始化,会 fence 掉旧实例(旧 epoch 的请求被拒绝),这是防止僵尸生产者重复写的关键;- 提交分两阶段:协调器先写 PrepareCommit 到
__transaction_state,再向相关分区写入控制消息(COMMIT/ABORT marker),最后完成; - 消费者用
isolation.level=read_committed,只读已提交事务的消息。
14.5 LSO 与 read_committed
事务消息写入后并非立即可见。Broker 为每个分区维护 LSO(Last Stable Offset):
日志: [普通消息] [事务T1消息] [普通消息] [事务T2消息] ...
^ ^
LSO 在这里 T2 未提交
read_committed 消费者最多读到 LSO。事务提交前,T1 的消息不可见;一旦提交,marker 之后的 LSO 前移,消费者才能读到。
监控事务相关指标:
transaction-coordinator-metrics下的active-transaction-count、prepare-transaction-completion-time;- 消费端
records-lag突增且isolation-level=read_committed时,检查是否有长时间未提交/未中止的事务。
14.6 Consume-Transform-Produce:流式精确一次
最常见的 exactly-once 形态:消费上游 topic,处理后写下游 topic,把“输入位移提交”和“输出写入”放进同一个事务。
props.put(ConsumerConfig.ISOLATION_LEVEL_CONFIG, "read_committed");
props.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, "etl-job-1");
consumer.subscribe(List.of("input-topic"));
producer.initTransactions();
while (true) {
var records = consumer.poll(Duration.ofMillis(500));
if (records.isEmpty()) continue;
producer.beginTransaction();
try {
for (var r : records) {
producer.send(new ProducerRecord<>(
"output-topic", r.key(), transform(r.value())));
}
// 把消费位移也纳入事务
Map<TopicPartition, OffsetAndMetadata> offsets = buildOffsets(records);
producer.sendOffsetsToTransaction(offsets, consumer.groupMetadata());
producer.commitTransaction();
} catch (Exception e) {
producer.abortTransaction();
// 重新 poll 会从上一次提交位移开始
}
}
这段代码保证了:输出消息与位移提交原子生效。要么下游看到结果且位移前移,要么两者都回滚,重放后再次写入。配合输入输出都在 Kafka,就构成了 Kafka 端到端的 exactly-once。
Kafka Streams 的 exactly_once_v2 本质就是把这套机制框架化,业务代码无需手写事务边界。
14.7 消费端的精确一次
如果下游是外部系统,事务管不到那里,需要自己设计幂等:
数据库方案
-- messages 表主键 = topic + partition + offset 或业务 request_id
INSERT INTO consumed_messages(message_key, topic_name, partition_id, offset_id)
VALUES (?, ?, ?, ?)
ON CONFLICT DO NOTHING;
插入成功才处理业务;插入冲突说明已处理过,直接跳过。把业务写和去重表放同一个本地事务,就是完整的端到端精确一次。
状态机方案
CREATED -> PAID -> SHIPPED
重复的 PAID 事件到来时,当前状态已是 SHIPPED -> 直接忽略
代价与选择
事务会带来吞吐下降(两阶段提交与 marker 写入)和延迟上升。问自己:
- 重复会造成资金/库存错误吗?必须严格幂等;
- 重复只是多记一条日志/多推一次通知吗?至少一次 + 去重可能就够;
- 输入输出是否都在 Kafka 内?是则优先用事务/Streams。
14.8 常见误区
- 开了幂等就不丢消息:幂等只防重,不丢消息靠
acks=all+ 副本配置; - 事务保证外部系统原子性:不能。MySQL 写入失败不会被 Kafka 事务回滚;
- read_committed 能去重业务消息:它只过滤未提交/已中止事务,不处理业务层重试造成的逻辑重复;
- 换个 transactional.id 重启:必须保持稳定,否则失去 fence 能力;
- 把所有业务都上事务:事务有真实成本,大多数场景“至少一次 + 幂等”更划算。
本章小结
- 至少一次 + 业务幂等是大多数系统的最佳平衡点;
- 幂等生产者用 PID+sequence 在 broker 端去重重试,但仅限单会话单分区;
- 事务提供跨分区原子写与 consume-transform-produce 的位移原子提交;
- 消费端用
read_committed只读已提交数据,边界由 LSO 控制; - 涉及外部系统时,精确一次要靠“去重表/状态机 + 本地事务”自己完成。
思考题
- 为什么幂等生产者不能防止“进程重启后重发同一业务消息”?该用什么解决?
sendOffsetsToTransaction如果漏掉,会发生什么?- 设计一个“Kafka -> Redis -> MySQL”链路的端到端精确一次方案,指出每一段由谁保证。