Producer 幂等与事务
定义与作用
Kafka Producer 的幂等和事务机制是两个递进的可靠性保证:
- 幂等生产者(Idempotent Producer):保证单分区内的消息不重复(Exactly-once 的第一层)
- 事务生产者(Transactional Producer):保证跨分区的消息原子性写入——要么全部成功,要么全部失败(Exactly-once 的第二层)
在分布式系统中,这两个机制解决了"网络不可靠"带来的根本矛盾:发送方不知道消息是否被 Broker 成功持久化,重试可能造成重复。
核心原理
幂等生产者:消除重试引起的重复
工作机制:
- Broker 为每个幂等 Producer 分配全局唯一的 Producer ID (PID)
- Producer 为每条消息附加单调递增的 Sequence Number (Seq)
- Broker 在内存中维护每个 Partition 的 (PID, Seq) 映射表,记录最近的 5 个 Seq
- 收到重复 Seq 的消息时,直接丢弃并返回成功
限制:
- 幂等只保证同一 Partition 内的不重复,跨分区不保证
- PID 在 Producer 重启后会变化——重启后的消息被视为新 Producer 的消息
事务生产者:原子性写入多分区
两阶段提交简化版:
- 所有参与分区的消息写入日志(标记为"未提交")
- Transaction Coordinator 写入 COMMIT Marker,消息变为"已提交"(对
read_committed消费者可见)
如果事务中止(abortTransaction()),写入 ABORT Marker,read_committed 消费者自动跳过这些消息。
完整示例
示例一:幂等生产者防止重复
场景:订单支付确认消息,由于网络超时被重试了 3 次。
操作前环境:Topic payments,单分区,Producer 配置幂等。
Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
// 开启幂等(自动设置 acks=all, retries=MAX, max.in.flight=5)
props.put("enable.idempotence", true);
try (KafkaProducer<String, String> producer = new KafkaProducer<>(props)) {
for (int i = 0; i < 100; i++) {
producer.send(new ProducerRecord<>("payments", "key-" + i, "payment-data-" + i));
}
}
操作后验证:
# 检查消息数 —— 应该是 100 条而非更多
bin/kafka-run-class.sh kafka.tools.GetOffsetShell \
--topic payments --time -1 --bootstrap-server localhost:9092
# payments:0:100 ← 正好 100 条,无重复
操作前后对比:
| 场景 | 不开幂等 | 开启幂等 |
|---|---|---|
| 正常发送 100 条 | 100 条 | 100 条 |
| 网络超时重试 10 次 | 100 ~ 110 条(有重复) | 100 条(幂等去重) |
示例二:事务生产者实现 Exactly-Once 跨 Topic 复制
场景:将 source-topic 的消息消费后写入 target-topic,同时提交 Consumer Offset,保证 Exactly-Once。
Properties consumerProps = new Properties();
consumerProps.put("bootstrap.servers", "localhost:9092");
consumerProps.put("group.id", "tx-copier");
consumerProps.put("enable.auto.commit", "false");
consumerProps.put("isolation.level", "read_committed");
consumerProps.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
consumerProps.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
Properties producerProps = new Properties();
producerProps.put("bootstrap.servers", "localhost:9092");
producerProps.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
producerProps.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
producerProps.put("transactional.id", "tx-copier-001"); // 事务 ID
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(consumerProps);
KafkaProducer<String, String> producer = new KafkaProducer<>(producerProps);
consumer.subscribe(Collections.singletonList("source-topic"));
producer.initTransactions();
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
if (records.isEmpty()) continue;
producer.beginTransaction();
try {
for (ConsumerRecord<String, String> record : records) {
producer.send(new ProducerRecord<>("target-topic",
record.key(), record.value()));
}
// 将 Consumer Offset 纳入同一事务
Map<TopicPartition, OffsetAndMetadata> offsets = new HashMap<>();
for (TopicPartition tp : records.partitions()) {
long offset = records.records(tp).get(records.records(tp).size() - 1).offset() + 1;
offsets.put(tp, new OffsetAndMetadata(offset));
}
producer.sendOffsetsToTransaction(offsets, "tx-copier");
producer.commitTransaction();
} catch (Exception e) {
producer.abortTransaction();
// 事务中止后,Consumer 需手动 seek 到上次已提交的 Offset 重新消费
}
}
操作前后对比:
| 情况 | 不开事务 | 开启事务 |
|---|---|---|
| 正常消费+写入 | OK | OK |
| 写入 target 后、提交 Offset 前崩溃 | target 有数据,但 Offset 未提交 → 重启后重复消费 → target 有重复 | 事务未提交 → target 和 Offset 都回滚 → 重启后重新处理,无重复 |
| 提交 Offset 后、写入 target 前崩溃 | Offset 已提交,但 target 无数据 → 丢消息 | 事务未提交 → 回滚 |
易错场景
易错 1:幂等生产者 + Producer 重启 = PID 变化
场景:依赖幂等去重,但 Producer 因 OOM 被 Kill 后自动重启(如 Kubernetes Pod 重启)。
后果:新的 Producer 实例获得新的 PID,旧 PID 的 Seq 去重失效。如果旧 Producer 有未确认的发送,新 Producer 可能产生重复。
正确做法:如果重启场景也需要 Exactly-Once 跨会话保证,使用事务生产者 + 稳定的 transactional.id。
易错 2:事务中忘记设置 isolation.level=read_committed
场景:Consumer 消费事务消息,但使用默认的 read_uncommitted。
后果:Consumer 会读到已中止事务的消息,破坏 Exactly-Once 语义。
解决:props.put("isolation.level", "read_committed");
面试高频考点
Q:Kafka 的 Exactly-Once 和数据库的 ACID 事务有什么不同?
A:
- 范围不同:Kafka 事务只在 Kafka 内部(Topic→Topic)提供 Exactly-Once;数据库事务覆盖表、行、索引
- 隔离级别:Kafka 只有
read_uncommitted和read_committed两种,不支持可重复读、串行化 - 实现方式:Kafka 通过 PID + Seq 去重 + COMMIT Marker 实现,而非锁和 MVCC
- 外部系统:要保证 Kafka → 外部系统(如 DB)的 Exactly-Once,需要外部系统配合(如两阶段提交、幂等写入)