乐途乐途
主页
  • 计算机基础

    • TCP/IP
    • Linux
    • HTTP
  • 数据库

    • SQL
    • MySQL 5.7
  • 编程语言

    • C
    • C++
    • Java SE
    • Python2
    • Python3
  • 数据格式

    • JSON
    • XML
  • 认证与安全

    • JWT
  • 工具

    • Markdown
  • Git

    • GitFlow
  • Quartz

    • Quartz
  • Java

    • Maven 入门
    • Maven 进阶
    • MyBatis
    • Spring
    • Spring MVC
  • Java

    • Spring Boot
    • Spring Cloud
    • Spring Cloud Alibaba
    • Spring Security
    • Spring AI
    • Spring Batch
    • Kafka
    • Java 设计模式
  • 缓存

    • Redis
  • 搜索引擎

    • Elasticsearch
  • 分布式协调

    • ZooKeeper
联系
阿里云
主页
  • 计算机基础

    • TCP/IP
    • Linux
    • HTTP
  • 数据库

    • SQL
    • MySQL 5.7
  • 编程语言

    • C
    • C++
    • Java SE
    • Python2
    • Python3
  • 数据格式

    • JSON
    • XML
  • 认证与安全

    • JWT
  • 工具

    • Markdown
  • Git

    • GitFlow
  • Quartz

    • Quartz
  • Java

    • Maven 入门
    • Maven 进阶
    • MyBatis
    • Spring
    • Spring MVC
  • Java

    • Spring Boot
    • Spring Cloud
    • Spring Cloud Alibaba
    • Spring Security
    • Spring AI
    • Spring Batch
    • Kafka
    • Java 设计模式
  • 缓存

    • Redis
  • 搜索引擎

    • Elasticsearch
  • 分布式协调

    • ZooKeeper
联系
阿里云
  • 学习路径
  • 第1章 消息队列与 Kafka 概述

    • 章节导读
    • 消息队列基础
    • Kafka 概述
    • Kafka 为什么快
  • 第2章 快速上手:单机环境搭建

    • 章节导读
    • 环境准备与安装
    • 快速启动
    • Topic 管理
  • 第3章 核心概念:主题、分区与日志

    • 章节导读
    • Topic
    • Partition
    • Offset
    • Segment 与存储结构
  • 第4章 生产者详解

    • 章节导读
    • Producer 概述
    • Producer 发送机制
    • Producer 分区策略
    • Producer 幂等与事务
    • Producer 配置
  • 第5章 消费者与消费组

    • 章节导读
    • Consumer 概述
    • Consumer Group 消费组
    • Consumer 分区分配策略
    • Consumer Offset 提交
    • Consumer 多线程
    • Consumer 配置
  • 第6章 Broker 与控制器

    • 章节导读
    • Broker 概述
    • Broker 配置
    • Controller
    • KRaft 共识协议
  • 第7章 副本与数据可靠性

    • 章节导读
    • 副本机制
    • ISR 与副本同步
    • Leader 选举
    • ACK 与一致性保证
    • 高水位与 Leader Epoch
  • 第8章 存储与性能优化

    • 章节导读
    • 存储架构
    • 日志清理与压缩
    • Page Cache 与零拷贝
    • Producer 性能优化
    • Consumer 性能优化
    • Broker 性能优化
  • 第9章 生产环境运维与监控

    • 章节导读
    • Topic 管理
    • Kafka 运维工具
    • 监控
    • 常见故障排查
  • 第10章 Kafka生态与面试考点

    • 章节导读
    • Kafka 生态全景
    • 面试高频 30 题

Producer 幂等与事务

定义与作用

Kafka Producer 的幂等和事务机制是两个递进的可靠性保证:

  • 幂等生产者(Idempotent Producer):保证单分区内的消息不重复(Exactly-once 的第一层)
  • 事务生产者(Transactional Producer):保证跨分区的消息原子性写入——要么全部成功,要么全部失败(Exactly-once 的第二层)

在分布式系统中,这两个机制解决了"网络不可靠"带来的根本矛盾:发送方不知道消息是否被 Broker 成功持久化,重试可能造成重复。

核心原理

幂等生产者:消除重试引起的重复

工作机制:

  1. Broker 为每个幂等 Producer 分配全局唯一的 Producer ID (PID)
  2. Producer 为每条消息附加单调递增的 Sequence Number (Seq)
  3. Broker 在内存中维护每个 Partition 的 (PID, Seq) 映射表,记录最近的 5 个 Seq
  4. 收到重复 Seq 的消息时,直接丢弃并返回成功

限制:

  • 幂等只保证同一 Partition 内的不重复,跨分区不保证
  • PID 在 Producer 重启后会变化——重启后的消息被视为新 Producer 的消息

事务生产者:原子性写入多分区

两阶段提交简化版:

  1. 所有参与分区的消息写入日志(标记为"未提交")
  2. 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 重新消费
    }
}

操作前后对比:

情况不开事务开启事务
正常消费+写入OKOK
写入 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:

  1. 范围不同:Kafka 事务只在 Kafka 内部(Topic→Topic)提供 Exactly-Once;数据库事务覆盖表、行、索引
  2. 隔离级别:Kafka 只有 read_uncommitted 和 read_committed 两种,不支持可重复读、串行化
  3. 实现方式:Kafka 通过 PID + Seq 去重 + COMMIT Marker 实现,而非锁和 MVCC
  4. 外部系统:要保证 Kafka → 外部系统(如 DB)的 Exactly-Once,需要外部系统配合(如两阶段提交、幂等写入)
上一页
Producer 分区策略
下一页
Producer 配置