Java 基础体系 · 第 97/100 篇。示例统一以 Java 25 LTS 为语言和 JVM 基线;框架示例使用与其兼容的现代稳定版本。

Java Kafka:Producer、Consumer、分区、Offset、事务和再均衡

Kafka 是一个以日志为核心的分布式消息系统。Producer(生产者)把记录追加到主题(Topic)的分区(Partition)中,Consumer(消费者)按分区顺序读取记录,并通过 Offset(偏移量)表示读取进度;Consumer Group(消费者组)负责把分区分配给组内消费者,再均衡(Rebalance)则负责在成员或分区发生变化时重新计算这种分配关系。

理解 Kafka 不能只记住几个 API。必须先明确一条数据从生产到消费的完整路径:

flowchart LR
    P[Producer] -->|序列化后的 Record| T[Topic]
    T --> A[Partition 0]
    T --> B[Partition 1]
    T --> C[Partition 2]

    A --> C1[Consumer A]
    B --> C2[Consumer B]
    C --> C1

    C1 -->|提交 offset| G[Consumer Group Coordinator]
    C2 -->|提交 offset| G
    G --> O[(__consumer_offsets)]

Producer 决定记录写入哪个分区;Broker 将记录追加到分区日志;Consumer 从指定分区拉取记录;消费者组协调器(Group Coordinator)保存每个消费者组在每个分区上的已提交 Offset。


一、Kafka 中最基本的数据模型

1. Topic、Partition 和 Record

Topic 是逻辑上的消息分类,例如:

orders
payments
user-events

Topic 通常由一个或多个 Partition 组成。每个分区是一个追加式日志:

Partition 0:
offset 0 -> record A
offset 1 -> record B
offset 2 -> record C

Offset 是分区内部的单调递增位置,不是 Topic 全局位置。因此下面两个记录可以同时存在:

orders-0 offset 10
orders-1 offset 10

它们属于不同分区,不能比较谁“更早”。

Kafka 中一条记录通常包含:

key
value
headers
timestamp
partition
offset

其中 partitionoffset 由 Kafka 最终确定,Producer 发送前可以指定 partition,也可以让分区器根据 key 或其他策略选择。

2. Kafka 的顺序保证范围

Kafka 的顺序保证是:

同一个分区内,记录按照追加顺序排列;消费者从该分区读取时,看到的是这个顺序。

Kafka 不保证同一个 Topic 的多个分区之间存在全局顺序。

例如:

orders-0: A(offset=0), C(offset=1)
orders-1: B(offset=0), D(offset=1)

消费者可能观察到:

A, B, C, D

也可能观察到:

B, A, D, C

如果业务要求同一个订单的事件严格有序,应当使用 orderId 作为 key,使同一个 orderId 的记录稳定进入同一个分区。


二、Producer:Java 如何生产记录

1. Producer 的核心职责

Java Kafka Producer 主要完成四件事:

  1. 将 Java 对象序列化为字节。
  2. 根据 Topic、Key 和分区器选择目标分区。
  3. 批量发送记录,并处理网络重试。
  4. 在需要时提供幂等性和事务能力。

Kafka Java 客户端的基础依赖是 kafka-clients。下面示例使用 Java 25 语法和标准 Kafka Producer API;Kafka 客户端版本应根据实际 Broker 版本和企业依赖管理策略确定。

<dependency>
    <groupId>org.apache.kafka</groupId>
    <artifactId>kafka-clients</artifactId>
    <version>${kafka.version}</version>
</dependency>

Java 25 本身不改变 Kafka 协议。Kafka 的 Producer、Consumer、事务和再均衡行为主要由 kafka-clients 版本以及 Broker 配置决定。

2. 最小 Producer 示例

import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.clients.producer.RecordMetadata;
import org.apache.kafka.common.serialization.StringSerializer;

import java.util.Properties;
import java.util.concurrent.Future;

public class SimpleProducer {
    public static void main(String[] args) throws Exception {
        Properties props = new Properties();
        props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
        props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
        props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());

        try (var producer = new KafkaProducer<String, String>(props)) {
            var record = new ProducerRecord<>(
                    "orders",
                    "order-1001",
                    "{\"orderId\":\"order-1001\",\"amount\":99}"
            );

            Future<RecordMetadata> future = producer.send(record);
            RecordMetadata metadata = future.get();

            System.out.printf(
                    "topic=%s, partition=%d, offset=%d%n",
                    metadata.topic(),
                    metadata.partition(),
                    metadata.offset()
            );
        }
    }
}

运行前必须满足:

  • Kafka Broker 可通过 localhost:9092 访问;
  • Topic orders 已存在,或者 Broker 开启了自动创建 Topic;
  • 客户端有写入权限;
  • Producer 使用的序列化器与 Java 类型匹配。

send() 通常是异步的。它将记录放入 Producer 的缓冲区,然后由后台 I/O 线程发送。future.get() 会等待 Broker 返回确认,因此示例可以观察到最终的分区和 Offset,但会牺牲异步发送带来的吞吐优势。

生产代码通常使用回调:

producer.send(record, (metadata, exception) -> {
    if (exception != null) {
        exception.printStackTrace();
        return;
    }

    System.out.printf(
            "sent partition=%d offset=%d%n",
            metadata.partition(),
            metadata.offset()
    );
});

回调中的异常不能忽略。send() 返回成功只代表记录被提交到客户端缓冲区,不代表 Broker 已经持久化成功。

3. Key 如何影响分区

当 ProducerRecord 没有显式指定分区时,分区器通常根据 key 选择分区。对于具有稳定 key 的记录,目标通常可以抽象为:

partition=hash(key)modNpartition = hash(key) \bmod N

其中:

  • key 是序列化前的逻辑键;
  • hash(key) 是分区器使用的哈希结果;
  • N 是当前 Topic 的分区数。

实际实现可能使用具体的哈希算法和分区器策略,因此不应把某个哈希细节当成跨版本协议保证。

当 key 相同且分区数不变时,记录通常进入同一个分区:

new ProducerRecord<>("orders", "order-1001", "CREATED");
new ProducerRecord<>("orders", "order-1001", "PAID");
new ProducerRecord<>("orders", "order-1001", "SHIPPED");

这样可以保证:

CREATED -> PAID -> SHIPPED

在同一分区中按追加顺序排列。

但是,增加分区数会改变分区映射。例如原来有 3 个分区:

hash(order-1001) % 3 = 1

增加到 4 个分区后可能变成:

hash(order-1001) % 4 = 3

因此,扩容分区可能破坏“同一 key 永远进入同一分区”的历史映射关系。Kafka 不会自动重排旧记录,新旧记录仍然保留在原分区和新分区中。

4. Producer 的确认、重试和幂等性

Producer 发送失败时可能重试。问题在于:Broker 可能已经写入记录,但响应在网络中丢失,Producer 无法确定写入是否成功。

如果 Producer 直接重试,可能出现:

第一次发送:Broker 已写入,但响应丢失
第二次发送:Broker 再写入一条相同记录
结果:消费者看到两条记录

幂等 Producer 通过 Producer ID、Epoch 和序列号让 Broker 识别重复发送,避免同一 Producer 会话中的重试造成重复追加。

现代 Kafka 客户端通常默认启用幂等相关行为,但生产系统仍应显式确认客户端版本和配置,不要仅依赖“默认值”。常见相关配置包括:

enable.idempotence=true
acks=all

acks=all 表示 Producer 等待 ISR(In-Sync Replicas,同步副本集合)确认。它提高持久性,但延迟和可用性会受到副本状态影响。

必须区分:

  • Producer 幂等性:避免 Producer 重试导致的重复写入;
  • 业务幂等性:消费者重复处理时,业务结果仍然正确;
  • 事务:将多个 Kafka 写入以及消费位点提交组合成原子操作。

Producer 幂等性不能让数据库更新自动幂等,也不能替代消费者的业务去重。


三、Consumer:拉取、处理和提交 Offset

1. Consumer 不是被 Broker 推送消息

Kafka Consumer 使用拉取模型。应用调用 poll(),向 Broker 请求当前分配分区上的记录。

ConsumerRecords<K, V> records = consumer.poll(Duration.ofMillis(1000));

poll() 返回的是一个批次,可能包含多个分区的记录。批次内的记录应按分区分别理解,不能把整个 ConsumerRecords 视为一个全局有序列表。

2. 最小 Consumer 示例

import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.common.serialization.StringDeserializer;

import java.time.Duration;
import java.util.List;
import java.util.Properties;

public class SimpleConsumer {
    public static void main(String[] args) {
        Properties props = new Properties();
        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
        props.put(ConsumerConfig.GROUP_ID_CONFIG, "order-service");
        props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG,
                StringDeserializer.class.getName());
        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG,
                StringDeserializer.class.getName());

        props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
        props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false");

        try (var consumer = new KafkaConsumer<String, String>(props)) {
            consumer.subscribe(List.of("orders"));

            while (true) {
                var records = consumer.poll(Duration.ofSeconds(1));

                for (ConsumerRecord<String, String> record : records) {
                    System.out.printf(
                            "topic=%s partition=%d offset=%d key=%s value=%s%n",
                            record.topic(),
                            record.partition(),
                            record.offset(),
                            record.key(),
                            record.value()
                    );
                }

                consumer.commitSync();
            }
        }
    }
}

这里的关键配置是:

  • group.id:消费者组标识;
  • enable.auto.commit=false:应用自己决定何时提交 Offset;
  • auto.offset.reset=earliest:没有已提交 Offset 时,从最早可用记录开始读取。

auto.offset.reset 只在“没有有效已提交 Offset”时生效。它不会覆盖已经存在的提交位置。常见值包括:

  • earliest:从分区最早可用位置读取;
  • latest:从当前末尾附近开始读取;
  • none:没有有效 Offset 时抛出异常。

3. Consumer 的生命周期

一个典型消费者的生命周期是:

创建 Consumer
    |
subscribe()
    |
加入消费者组
    |
分配分区
    |
poll()
    |
处理记录
    |
提交 Offset
    |
继续 poll()
    |
关闭或触发再均衡

subscribe() 并不会立即完成分区分配。真正的分配通常发生在后续 poll() 过程中,因为 Consumer 需要与 Group Coordinator 通信并参加消费者组协议。

Consumer 必须持续调用 poll()。即使暂时没有业务记录,也需要调用 poll() 来维持组成员身份、执行心跳相关工作并接收再均衡结果。

如果一次批次处理时间过长,超过 max.poll.interval.ms,Broker 可能认为该消费者失效并触发再均衡。此时处理线程可能仍在执行旧批次,但该消费者已经失去分区所有权。

4. Offset 到底表示什么

对于分区:

offset 0 -> A
offset 1 -> B
offset 2 -> C
offset 3 -> 尚未写入

如果消费者处理完 AB,下一条要读取的是 C,那么通常提交的是:

offset = 2

也就是说,Consumer 提交的 Offset 通常表示:

下一次要读取的位置,而不是最后一条已处理记录的位置。

Java API 中常见写法:

consumer.commitSync();

或者显式提交某个分区的下一位置:

import org.apache.kafka.clients.consumer.OffsetAndMetadata;
import org.apache.kafka.common.TopicPartition;

var offsets = Map.of(
        new TopicPartition("orders", 0),
        new OffsetAndMetadata(2L)
);

consumer.commitSync(offsets);

如果提交 2,表示分区 orders-0 的 offset 01 已经被认为处理完成,恢复后从 2 开始。

5. 提交时机决定消息语义

假设一批记录为:

offset 10 -> A
offset 11 -> B
offset 12 -> C

先处理,再提交

处理 A
处理 B
处理 C
提交 13

如果进程在提交前崩溃,重启后会重新读取 10、11、12。因此可能重复处理,但不会因为提前提交而跳过尚未处理的记录。

这种方式通常对应 至少一次(At-Least-Once) 处理语义:

不丢失作为优先目标,但允许重复

先提交,再处理

提交 13
处理 A
处理 B
处理 C

如果提交成功后进程在处理 B 时崩溃,重启会从 13 开始,BC 不再被读取。于是可能丢失尚未完成的业务处理。

处理一条就提交一条

处理 A -> 提交 11
处理 B -> 提交 12
处理 C -> 提交 13

重复窗口更小,但提交请求更多,吞吐和延迟会受到影响。

commitAsync() 延迟更低,但回调失败处理复杂;commitSync() 会等待结果,错误更容易处理,但可能阻塞消费线程。无论使用哪种提交方式,都不能把“提交成功”误认为“外部业务一定成功”,除非提交与业务操作具有同一事务边界。


四、Consumer Group 和分区分配

1. 一个分区同一时间只属于组内一个消费者

对于同一个消费者组:

分区数 = 3
消费者数 = 2

一种分配结果可能是:

Consumer A -> Partition 0, Partition 1
Consumer B -> Partition 2

如果消费者数大于分区数:

分区数 = 2
消费者数 = 3

则一定有一个消费者没有分区:

Consumer A -> Partition 0
Consumer B -> Partition 1
Consumer C -> 无分区

因此,增加消费者数量不能突破 Topic 分区数提供的并行度上限。

对于一个分区,同一个消费者组内不会同时由两个消费者处理。但是不同消费者组可以独立读取同一个 Topic:

order-service-group -> orders
audit-service-group  -> orders

这两个组各自维护 Offset,互不影响。

2. 分配策略和实际限制

Kafka 会使用分区分配策略在组内分配分区。常见策略包括 Range、RoundRobin、Sticky 和 CooperativeSticky。具体可用策略取决于客户端版本和配置。

分配策略解决的是“哪些分区交给哪些消费者”,但不解决以下问题:

  • 业务处理是否幂等;
  • 外部数据库是否提交成功;
  • 一个批次内部如何并行处理;
  • 处理线程是否会阻塞 poll()
  • 消费者是否有权限读取 Topic。

如果一个 Consumer 从多个分区读取记录并使用线程池并行处理,必须额外维护每个分区的连续完成位置。不能因为 offset 12 完成,就提交 offset 13;如果 offset 11 仍在处理中,恢复后从 13 开始会跳过 11。


五、再均衡:为什么发生、发生了什么

1. 再均衡的定义

再均衡(Rebalance) 是消费者组重新计算分区归属的过程。它通常发生在以下情况:

  • 新消费者加入消费者组;
  • 消费者正常退出;
  • 消费者崩溃或心跳超时;
  • Topic 分区数量变化;
  • 消费者订阅的 Topic 集合变化;
  • 组协议或分配策略发生变化。

再均衡的本质不是“重新读取消息”,而是:

旧的分区所有权失效
    ->
消费者组重新协商成员
    ->
新的分区分配产生
    ->
消费者从各自已提交位置继续读取

2. 典型再均衡过程

sequenceDiagram
    participant A as Consumer A
    participant B as Consumer B
    participant G as Group Coordinator
    participant K as Kafka Brokers

    A->>G: 加入消费者组
    B->>G: 加入消费者组
    G->>K: 获取 Topic 分区信息
    G-->>A: 分配 P0、P1
    G-->>B: 分配 P2

    A->>A: 处理 P0、P1
    B->>B: 处理 P2

    Note over A: A 崩溃或 poll 超时
    G->>G: 检测成员失效
    G->>G: 触发再均衡
    G-->>B: 分配 P0、P1、P2
    B->>K: 从已提交 Offset 继续读取

在再均衡期间,分区所有权会发生变化。旧消费者如果继续使用已经被撤销的分区,可能产生提交失败、重复处理或业务并发冲突。

3. ConsumerRebalanceListener

需要在再均衡前保存处理进度时,可以使用监听器:

import org.apache.kafka.clients.consumer.ConsumerRebalanceListener;
import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.common.TopicPartition;

import java.util.Collection;
import java.util.Map;

consumer.subscribe(
        List.of("orders"),
        new ConsumerRebalanceListener() {
            @Override
            public void onPartitionsRevoked(
                    Collection<TopicPartition> partitions) {

                // 停止向这些分区提交新的业务处理结果
                // 刷新已经完成的处理进度
                consumer.commitSync();
            }

            @Override
            public void onPartitionsAssigned(
                    Collection<TopicPartition> partitions) {

                // 新分区已经分配,可以初始化分区相关状态
                System.out.println("assigned: " + partitions);
            }
        }
);

上例展示了生命周期,但实际使用时需要谨慎:

  • onPartitionsRevoked 中不要执行不可控的长时间阻塞操作;
  • 如果使用协作式再均衡,撤销回调接收到的分区语义与立即撤销策略不同;
  • 不能无条件 commitSync() 就认为所有异步任务都已完成;
  • 需要区分“已拉取”“已开始处理”和“已完成处理”。

4. 再均衡导致重复的具体路径

假设:

P0:
offset 100 -> A
offset 101 -> B
offset 102 -> C

消费者处理完 A,但尚未提交;处理 B 时发生再均衡,P0 被交给另一个消费者。

如果上次提交位置仍是:

offset = 100

新消费者会从 100 开始读取:

A、B、C

于是 A 和已经处理过的 B 可能重复。

这不是 Kafka “重复发送”造成的,而是“业务处理进度”和“已提交消费位置”之间存在时间差。只要采用至少一次语义,就必须设计幂等处理。

5. 处理过慢会主动触发再均衡

Consumer 有两个容易混淆的时间概念:

  • session.timeout.ms:协调器多久没有收到有效心跳后认为成员失效;
  • max.poll.interval.ms:两次 poll() 之间允许的最大间隔,超过后认为消费者处理能力不足或失去响应。

例如:

poll() 返回 5000 条记录
业务处理耗时 10 分钟
max.poll.interval.ms = 5 分钟

即使底层心跳仍可能存在,超过最大 poll 间隔后也可能触发消费者离组。

解决路径不是盲目增大超时时间,而是先判断处理模型:

  1. 减小 max.poll.records
  2. 将长任务拆分为更小的处理单元;
  3. 使用受控的工作线程池;
  4. 确保提交位置只推进到连续完成的位置;
  5. 对无法及时处理的消息设计重试或死信流程;
  6. 必要时调整 max.poll.interval.ms,但要接受故障检测变慢的代价。

六、分区、并行度和吞吐的推导

假设一个 Topic 有 P 个分区,一个消费者组有 C 个消费者,每个消费者只有一个消费线程,则有效消费并行度近似为:

parallelismmin(P,C)parallelism \leq \min(P, C)

如果每个消费者内部再使用线程池,理论上可以启动更多处理任务,但分区内顺序约束仍然存在。

对于单分区,如果要求严格顺序:

offset 10 -> offset 11 -> offset 12

那么 offset 11 不能在 offset 10 之前提交,即使 offset 11 的业务计算先完成。可提交位置必须是“从当前起点开始,连续完成的最大前缀”。

例如:

处理完成状态:
10 -> 完成
11 -> 未完成
12 -> 完成

虽然 12 完成了,也只能提交:

11

因为 11 尚未完成,提交 13 会跳过它。

如果每个分区每秒最多处理 R 条记录,P 个分区的理论处理能力约为:

throughputP×Rthroughput \leq P \times R

这是简化模型。实际吞吐还受到以下因素影响:

  • Producer 批量大小;
  • Broker 磁盘和网络;
  • 副本同步;
  • Consumer 拉取批次;
  • 业务处理延迟;
  • 外部数据库;
  • 再均衡和重试;
  • 单条消息大小。

增加分区可以提高并行度,但也会增加文件、网络连接、元数据和再均衡管理成本,并且可能改变 key 到分区的映射。


七、Offset 的三种位置和常见误解

Kafka Consumer 至少应区分三个位置:

1. Log Start Offset

分区当前仍然保留的最早 Offset。由于日志保留策略,旧记录可能已经被删除。

2. Position

当前 Consumer 下一次要读取的位置。它是 Consumer 本地运行状态。

3. Committed Offset

消费者组已经提交到 Kafka 的位置。它是故障恢复时使用的持久化进度。

示例:

Log Start Offset = 100
Committed Offset = 120
Position         = 125

这表示:

  • 100 之前的数据已经被删除或不可读取;
  • 组恢复时从 120 开始;
  • 当前 Consumer 已经拉取到 125;
  • 120 到 124 的处理结果可能尚未提交。

因此,看到 Consumer 当前 position 前进,并不表示业务处理已经安全完成。

4. Offset 提交失败的处理

如果 commitSync() 抛出异常,不应继续假设提交成功。可以:

  • 停止继续推进业务状态;
  • 记录失败的分区和 Offset;
  • 让 Consumer 重新处理未确认部分;
  • 根据异常类型判断是否需要退出或等待恢复。

在 Consumer 已经失去分区所有权后提交,可能抛出 CommitFailedException。这通常说明处理时间过长或发生了再均衡,而不是简单的网络重试问题。


八、事务:把 Kafka 内的多步操作绑定起来

1. Kafka 事务解决什么问题

假设一个应用:

  1. orders 读取订单;
  2. 计算库存事件;
  3. inventory-events 发送结果;
  4. 提交消费 Offset。

如果第 3 步成功、第 4 步失败,应用重启后会再次读取订单并再次发送库存事件。

如果第 4 步成功、第 3 步实际失败或结果不可见,则可能出现数据丢失。

Kafka 事务可以把以下操作放在同一个原子边界内:

消费输入记录
    +
向一个或多个 Kafka Topic 生产输出
    +
提交消费者组 Offset

事务提交成功时,输出记录和消费进度一起生效;事务中止时,输出记录不会对 read_committed 消费者可见,消费 Offset 也不会按事务提交。

2. Producer 事务的基本代码

import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.common.serialization.StringSerializer;

import java.util.Properties;

public class TransactionalProducer {
    public static void main(String[] args) {
        Properties props = new Properties();
        props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
        props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG,
                StringSerializer.class.getName());
        props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG,
                StringSerializer.class.getName());

        props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, "true");
        props.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG,
                "order-producer-instance-1");

        try (var producer = new KafkaProducer<String, String>(props)) {
            producer.initTransactions();

            try {
                producer.beginTransaction();

                producer.send(new ProducerRecord<>(
                        "order-events",
                        "order-1001",
                        "ORDER_CREATED"
                ));

                producer.send(new ProducerRecord<>(
                        "audit-events",
                        "order-1001",
                        "ORDER_CREATED_AUDIT"
                ));

                producer.commitTransaction();
            } catch (RuntimeException e) {
                producer.abortTransaction();
                throw e;
            }
        }
    }
}

transactional.id 是事务 Producer 的稳定标识。它用于让 Broker 识别同一个事务 Producer 的新旧实例,并隔离旧实例继续写入的能力。

同一个 transactional.id 不应被多个仍然活跃的 Producer 实例同时使用。例如两个应用实例错误地配置成相同的 ID,较新的 Producer 可能使旧 Producer 失效,旧 Producer 后续发送或提交事务时抛出异常。

3. 事务的状态变化

一个事务可以抽象为:

未开始
  |
beginTransaction()
  |
进行中
  | \
  |  \ abortTransaction()
  |   \
  |   已中止
  |
commitTransaction()
  |
已提交

如果发送失败、事务超时或 Producer 被 Broker 判定为过期,事务不能继续正常提交,应用必须根据异常进行中止、重试或重启 Producer。

事务不是普通的 try/catch。下面的写法不能形成事务:

try {
    producer.send(record1);
    producer.send(record2);
} catch (Exception e) {
    // 这里没有撤销已发送记录的能力
}

只有显式调用 beginTransaction()commitTransaction()abortTransaction(),并且 Broker 支持事务相关配置时,才有 Kafka 事务语义。


九、消费—处理—生产的事务性工作流

Kafka 最典型的事务场景是:

输入 Topic -> Consumer -> 业务转换 -> 输出 Topic

Java Consumer 需要取得消费者组元数据,然后把 Offset 作为事务的一部分提交给 Producer。

import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.OffsetAndMetadata;
import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.common.TopicPartition;
import org.apache.kafka.common.serialization.StringDeserializer;
import org.apache.kafka.common.serialization.StringSerializer;

import java.time.Duration;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Properties;

public class ConsumeTransformProduce {
    public static void main(String[] args) {
        var producerProps = new Properties();
        producerProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
        producerProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG,
                StringSerializer.class.getName());
        producerProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG,
                StringSerializer.class.getName());
        producerProps.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG,
                "transformer-instance-1");

        var consumerProps = new Properties();
        consumerProps.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
        consumerProps.put(ConsumerConfig.GROUP_ID_CONFIG, "transformer-group");
        consumerProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG,
                StringDeserializer.class.getName());
        consumerProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG,
                StringDeserializer.class.getName());
        consumerProps.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false");
        consumerProps.put(ConsumerConfig.ISOLATION_LEVEL_CONFIG, "read_committed");

        try (var producer = new KafkaProducer<String, String>(producerProps);
             var consumer = new org.apache.kafka.clients.consumer.KafkaConsumer
                     <String, String>(consumerProps)) {

            producer.initTransactions();
            consumer.subscribe(List.of("order-events"));

            while (true) {
                var records = consumer.poll(Duration.ofSeconds(1));
                if (records.isEmpty()) {
                    continue;
                }

                producer.beginTransaction();

                try {
                    for (ConsumerRecord<String, String> record : records) {
                        String output = transform(record.value());

                        producer.send(new ProducerRecord<>(
                                "normalized-orders",
                                record.key(),
                                output
                        ));
                    }

                    Map<TopicPartition, OffsetAndMetadata> offsets =
                            new HashMap<>();

                    records.partitions().forEach(partition -> {
                        var partitionRecords = records.records(partition);
                        long nextOffset =
                                partitionRecords.getLast().offset() + 1;

                        offsets.put(
                                partition,
                                new OffsetAndMetadata(nextOffset)
                        );
                    });

                    producer.sendOffsetsToTransaction(
                            offsets,
                            consumer.groupMetadata()
                    );

                    producer.commitTransaction();
                } catch (RuntimeException e) {
                    producer.abortTransaction();
                    throw e;
                }
            }
        }
    }

    private static String transform(String input) {
        return input.toUpperCase();
    }
}

这段流程的关键顺序是:

  1. Consumer poll() 获取输入记录;
  2. Producer 开始事务;
  3. 发送输出记录;
  4. 计算每个输入分区的下一 Offset;
  5. 调用 sendOffsetsToTransaction()
  6. 提交事务。

只有第 6 步成功,输出记录和 Offset 才一起生效。

records 可能包含多个分区,因此必须按分区计算提交位置。对于每个分区,提交值是该批次最后一条记录 Offset 加一,而不是所有分区共用一个最大 Offset。

事务性消费的前提

事务性消费通常还需要:

enable.auto.commit=false
isolation.level=read_committed

Producer 侧需要:

enable.idempotence=true
transactional.id=稳定且唯一的实例标识

Broker 侧需要正确配置事务日志和副本相关参数。事务依赖 Kafka 集群内部的事务协调机制;如果事务状态日志不可用,Producer 不能正常完成事务生命周期。

事务不能覆盖外部系统

下面的流程仍然不是一个跨系统原子事务:

Kafka Consumer
    -> 更新 MySQL
    -> Kafka Producer 提交事务

如果 MySQL 已提交,但 Kafka 事务提交失败,数据库和 Kafka 仍然不一致。Kafka 事务只覆盖 Kafka 参与的操作,不能自动加入任意数据库、HTTP 服务或文件系统。

要解决 Kafka 与数据库的一致性问题,通常需要:

  • Outbox Pattern;
  • 数据库本地事务记录事件,再由 CDC 发布;
  • 可重试且幂等的补偿流程;
  • 明确接受最终一致性。

不能因为调用了 sendOffsetsToTransaction(),就认为数据库更新也被纳入事务。


十、read_committedread_uncommitted

Kafka 事务写入的记录会经历事务状态变化。消费者读取时有两种主要隔离级别:

read_uncommitted

消费者可以读取尚未提交事务中的记录,也可以读取最终被中止事务写入的记录。

适合不关心事务可见性的场景,但可能看到业务上不应暴露的中间结果。

read_committed

消费者只读取已提交事务中的记录。中止事务写入的记录对它不可见。

isolation.level=read_committed

read_committed 不代表“业务一定成功”,它只表示 Kafka 事务已经提交。若事务内的业务逻辑本身计算错误,错误结果仍然可以被提交并被读取。

另外,事务记录和事务标记会影响消费者的读取边界。一个事务尚未结束时,read_committed Consumer 可能暂时无法继续越过该事务的 Last Stable Offset(稳定可见边界),因此长事务会增加消费延迟。


十一、至少一次、至多一次和恰好一次

1. 至多一次(At-Most-Once)

流程:

先提交 Offset
再处理业务

故障时可能丢失记录,但通常不会重复处理。

适用于允许少量丢失、希望降低重复成本的场景。它不是“没有失败”,而是把失败结果偏向丢失。

2. 至少一次(At-Least-Once)

流程:

先处理业务
再提交 Offset

故障时可能重新读取已处理记录,但只要 Offset 未被错误提前提交,就不容易跳过未处理数据。

这是 Kafka 消费应用中常见的默认取舍。业务处理必须具备幂等性,例如:

INSERT INTO processed_event(event_id, result)
VALUES (?, ?)
ON CONFLICT (event_id) DO NOTHING;

上面的 SQL 是 PostgreSQL 风格示例。它利用 event_id 唯一约束防止同一事件重复产生业务结果;其他数据库需要使用对应的唯一键和冲突处理语法。

3. 恰好一次(Exactly-Once)

在 Kafka 内部,Consume-Transform-Produce 可以借助 Kafka 事务实现更强的恰好一次处理语义:

输入 Offset 提交
+
输出 Topic 写入

但是“恰好一次”必须限定边界:

  • Kafka 到 Kafka:可以通过事务实现原子提交和读取可见性;
  • Kafka 到数据库:不能仅靠 Kafka 事务实现;
  • Kafka 到外部 HTTP:不能保证对方只执行一次;
  • 业务副作用:仍需要幂等键、去重或补偿。

所以更准确的表述是:

Kafka 事务可以提供 Kafka 读写链路内的 Exactly-Once Semantics,而不是整个现实世界业务流程的绝对只执行一次。


十二、失败路径分析

1. Producer 在发送后崩溃

Producer.send()
    -> Broker 已写入
    -> Producer 在收到响应前崩溃

重启后如果重新发送同一业务事件:

  • 非幂等 Producer 可能产生重复记录;
  • 幂等 Producer 可以处理同一 Producer 会话中的重试;
  • 新进程重新发送仍可能需要业务事件 ID 去重。

因此记录通常应包含稳定的业务标识:

{
  "eventId": "evt-2025-0001",
  "orderId": "order-1001",
  "type": "ORDER_PAID"
}

2. Consumer 处理成功但提交失败

处理业务成功
    ->
commitSync() 失败
    ->
进程继续或崩溃

恢复后会重新读取该记录。业务必须允许重复,或者通过 eventId 做去重。

3. Consumer 提交成功但业务处理失败

commitSync() 成功
    ->
业务处理失败

恢复后不会重新读取已经提交的记录,这会形成逻辑丢失。根因通常是提交时机错误,而不是 Kafka 读取错误。

4. 事务提交超时

事务提交超时不一定能直接推断“Broker 没有提交”。网络故障可能发生在提交已生效之后。应用不能简单地无条件开始一个新事务并重复写入而不考虑幂等性。

事务 Producer 遇到不可恢复状态时,通常需要关闭并重新创建 Producer;是否可以继续使用当前实例取决于具体异常类型。Kafka 客户端会区分可重试异常、事务异常和 Producer fencing 等状态。


十三、死信、重试与 Poison Message

Poison Message 指一条无论如何都无法被当前业务成功处理的记录,例如:

  • 数据格式永久错误;
  • 必填字段缺失;
  • 版本不兼容;
  • 业务约束永远不满足。

如果消费者对它无限重试:

poll -> 处理失败 -> 不提交 -> 重启或重试 -> 同一条记录再次失败

后续记录可能永远无法处理,因为同一分区必须按顺序推进。

一种常见流程是:

主 Topic
   |
   | 失败且达到重试上限
   v
重试 Topic / 延迟 Topic
   |
   | 仍然失败
   v
死信 Topic

死信记录应保留足够诊断信息:

{
  "originalTopic": "orders",
  "originalPartition": 0,
  "originalOffset": 123,
  "errorType": "ValidationException",
  "errorMessage": "amount is missing",
  "failedAt": "2025-01-01T00:00:00Z",
  "originalPayload": "..."
}

如果“发布到死信 Topic”和“提交原始 Topic Offset”需要原子完成,可以把它们放在同一个 Kafka 事务中。否则可能出现死信已发布但原始 Offset 未提交,导致死信重复;也可能 Offset 已提交但死信发布失败,导致原始记录丢失。


十四、使用命令验证分区和 Offset

假设 Kafka 地址为 localhost:9092,可以使用 Kafka 自带命令查看 Topic:

kafka-topics.sh \
  --bootstrap-server localhost:9092 \
  --describe \
  --topic orders

典型结果类似:

Topic: orders  PartitionCount: 3
Partition: 0  Leader: 1  Replicas: 1  Isr: 1
Partition: 1  Leader: 1  Replicas: 1  Isr: 1
Partition: 2  Leader: 1  Replicas: 1  Isr: 1

重点观察:

  • PartitionCount:分区数;
  • Leader:当前负责读写的副本;
  • Replicas:副本集合;
  • Isr:当前同步副本集合。

查看消费者组位点:

kafka-consumer-groups.sh \
  --bootstrap-server localhost:9092 \
  --describe \
  --group order-service

典型结果类似:

GROUP          TOPIC   PARTITION  CURRENT-OFFSET  LOG-END-OFFSET  LAG  CONSUMER-ID
order-service  orders  0          120             125             5    consumer-a
order-service  orders  1          98              98              0    consumer-b

这些字段表示:

  • CURRENT-OFFSET:消费者组已提交的 Offset;
  • LOG-END-OFFSET:分区当前末尾位置;
  • LAG:通常可理解为两者之差;
  • CONSUMER-ID:当前持有分区的消费者。

例如分区 0:

LOG-END-OFFSET - CURRENT-OFFSET = 125 - 120 = 5

表示组的提交位置落后于日志末尾 5 个位置。Lag 增大可能是消费处理变慢,也可能是 Producer 突然增速,不能单独据此判断故障。

Offset 重置风险

Kafka 提供消费者组 Offset 重置能力,但这是高风险操作。重置前必须确认:

  1. 当前消费者组已经停止,避免并发消费;
  2. 目标 Topic 和分区范围正确;
  3. 新 Offset 的时间或位置正确;
  4. 重置后会不会产生大量重复业务;
  5. 外部数据库是否支持重复处理;
  6. 是否已经记录当前 Offset,便于恢复。

错误地把 Offset 重置到 earliest,可能造成全量历史数据重新处理;错误地重置到 latest,可能跳过尚未处理的历史数据。


十五、Spring Boot 中的配置边界

如果使用 Spring Boot,底层仍然是 Kafka Java Client。Spring Kafka 会负责 Listener 容器、线程生命周期、错误处理和提交协调,但不会改变 Kafka 分区、Offset 和事务的基本语义。

一个基础配置可以写成:

spring:
  kafka:
    bootstrap-servers: localhost:9092
    consumer:
      group-id: order-service
      enable-auto-commit: false
      auto-offset-reset: earliest
      properties:
        isolation.level: read_committed
    producer:
      properties:
        enable.idempotence: true
        transactional.id: order-service-${INSTANCE_ID}

这里的 transactional.id 必须按实例隔离。不能让多个活跃实例共享同一个事务 ID,否则会发生 Producer fencing。

Spring 容器中的监听方法何时提交 Offset,取决于容器的 AckMode、异常处理器和事务配置。不能只看到:

@KafkaListener(topics = "orders")
public void listen(String value) {
    service.process(value);
}

就断定它一定是至少一次、恰好一次或自动提交。必须同时检查:

  • Consumer enable.auto.commit
  • Listener Container 的 AckMode;
  • 异常是否被吞掉;
  • 是否配置 KafkaTransactionManager;
  • 输出 Producer 是否加入同一 Kafka 事务;
  • 外部数据库是否参与同一事务。

Spring Boot 的自动配置简化了对象创建,但 Offset 提交和异常传播仍然决定实际消息语义。


十六、诊断问题时应沿着状态链路检查

遇到“消息没消费”“消息重复”“消息顺序错乱”时,应按以下因果链检查,而不是先修改超时参数。

消息没有被消费

依次确认:

Topic 是否存在
    ->
Consumer 是否订阅成功
    ->
group.id 是否正确
    ->
Consumer 是否分配到分区
    ->
已提交 Offset 是否已经超过目标记录
    ->
auto.offset.reset 是否真正生效
    ->
Consumer 是否有读取权限

如果消费者组已有提交位置,修改 auto.offset.reset=earliest 通常不会让它重新从头读取。

消息重复

检查:

业务处理是否在提交前完成
    ->
commit 是否失败
    ->
是否发生再均衡
    ->
是否有多个实例使用同一 transactional.id
    ->
Producer 是否因不确定响应而重试
    ->
业务是否缺少 eventId 幂等

重复不一定意味着 Kafka Broker 保存了重复消息,也可能只是同一条记录被 Consumer 从未提交的位置再次读取。

顺序错误

先判断“顺序”属于哪个范围:

  • 同一分区内是否乱序;
  • 同一个 key 是否被路由到多个分区;
  • Consumer 是否并行处理同一分区;
  • 业务是否异步完成后乱序提交;
  • 是否在增加分区后改变了 key 映射。

如果同一个业务实体跨多个分区发送,Kafka 无法提供跨分区顺序。需要重新设计 key 或在业务层使用版本号、序列号检测乱序。

再均衡频繁发生

重点观察:

两次 poll 的间隔
业务处理批次耗时
max.poll.interval.ms
心跳与 session timeout
消费者进程是否频繁重启
Topic 分区是否频繁变化

如果处理线程阻塞在数据库、远程 HTTP 或锁等待上,增大 session.timeout.ms 可能没有帮助,因为真正触发的是 max.poll.interval.ms 或成员离组。


十七、规范保证、实现行为和工程取舍

需要明确区分三类结论。

Kafka 协议和模型提供的保证

  • Offset 只在分区内有意义;
  • 同一分区中的记录具有追加顺序;
  • 同一消费者组内,一个分区同一时间只分配给一个消费者;
  • 消费者组 Offset 独立保存;
  • Kafka 事务可以原子提交 Kafka 输出和消费者组 Offset;
  • read_committed 不读取中止事务的记录。

依赖客户端或 Broker 版本的行为

  • 默认分区器的具体实现;
  • Producer 幂等性默认配置;
  • 可用的再均衡协议和策略;
  • 事务相关默认超时;
  • Spring Boot 和 Spring Kafka 的容器默认 AckMode。

这些内容不能脱离实际版本直接推断。Java 25 LTS 只是运行时版本,Kafka 客户端和 Spring Boot 仍需单独确认兼容性。

需要由业务设计承担的责任

  • 外部数据库的一致性;
  • 重复消费的幂等处理;
  • Poison Message 的隔离;
  • 跨分区顺序;
  • 远程调用失败后的补偿;
  • Offset 重置造成的重复或丢失影响。

Producer、Consumer、分区、Offset、事务和再均衡并不是互相独立的功能点:

分区决定并行度和顺序边界
    ->
Consumer Group 决定分区归属
    ->
再均衡改变分区所有权
    ->
Offset 决定故障恢复位置
    ->
提交时机决定重复或丢失风险
    ->
事务可以把 Kafka 输出与 Offset 绑定
    ->
外部系统仍需要独立的一致性方案

掌握这条状态链路后,Kafka 的多数问题都可以还原为:记录在哪个分区、Consumer 当前读到哪里、组提交了哪里、分区是否刚刚被撤销,以及业务副作用是否和 Offset 处于同一个原子边界。


系列导航与关联阅读

官方资料

本文依据 Java、Spring 与相关项目官方文档重新梳理;正文、示例与生产清单由 WR BLOG 编写。