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
其中 partition 和 offset 由 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 主要完成四件事:
- 将 Java 对象序列化为字节。
- 根据 Topic、Key 和分区器选择目标分区。
- 批量发送记录,并处理网络重试。
- 在需要时提供幂等性和事务能力。
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 的记录,目标通常可以抽象为:
其中:
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 -> 尚未写入
如果消费者处理完 A 和 B,下一条要读取的是 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 0 和 1 已经被认为处理完成,恢复后从 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 开始,B 和 C 不再被读取。于是可能丢失尚未完成的业务处理。
处理一条就提交一条
处理 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 间隔后也可能触发消费者离组。
解决路径不是盲目增大超时时间,而是先判断处理模型:
- 减小
max.poll.records; - 将长任务拆分为更小的处理单元;
- 使用受控的工作线程池;
- 确保提交位置只推进到连续完成的位置;
- 对无法及时处理的消息设计重试或死信流程;
- 必要时调整
max.poll.interval.ms,但要接受故障检测变慢的代价。
六、分区、并行度和吞吐的推导
假设一个 Topic 有 P 个分区,一个消费者组有 C 个消费者,每个消费者只有一个消费线程,则有效消费并行度近似为:
如果每个消费者内部再使用线程池,理论上可以启动更多处理任务,但分区内顺序约束仍然存在。
对于单分区,如果要求严格顺序:
offset 10 -> offset 11 -> offset 12
那么 offset 11 不能在 offset 10 之前提交,即使 offset 11 的业务计算先完成。可提交位置必须是“从当前起点开始,连续完成的最大前缀”。
例如:
处理完成状态:
10 -> 完成
11 -> 未完成
12 -> 完成
虽然 12 完成了,也只能提交:
11
因为 11 尚未完成,提交 13 会跳过它。
如果每个分区每秒最多处理 R 条记录,P 个分区的理论处理能力约为:
这是简化模型。实际吞吐还受到以下因素影响:
- 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 事务解决什么问题
假设一个应用:
- 从
orders读取订单; - 计算库存事件;
- 向
inventory-events发送结果; - 提交消费 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();
}
}
这段流程的关键顺序是:
- Consumer
poll()获取输入记录; - Producer 开始事务;
- 发送输出记录;
- 计算每个输入分区的下一 Offset;
- 调用
sendOffsetsToTransaction(); - 提交事务。
只有第 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_committed 与 read_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 重置能力,但这是高风险操作。重置前必须确认:
- 当前消费者组已经停止,避免并发消费;
- 目标 Topic 和分区范围正确;
- 新 Offset 的时间或位置正确;
- 重置后会不会产生大量重复业务;
- 外部数据库是否支持重复处理;
- 是否已经记录当前 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 完整学习路线:从 Java 25 语言与 JVM 到 Spring、微服务和生产交付
- 上一篇:Java gRPC:Protobuf、Unary、Stream、拦截器、Deadline 和治理
- 下一篇:Java RabbitMQ:Exchange、确认、重试、死信、顺序和幂等
官方资料
本文依据 Java、Spring 与相关项目官方文档重新梳理;正文、示例与生产清单由 WR BLOG 编写。

评论
0 条讨论