Chapter 03
上手实操:把原理写成可运行的 Java
第 2 章把那条日志拆开讲了机制——副本与 ISR 怎么保证持久性、投递语义的真实边界。理解了机制,这一章把它落成代码:用 kafka-clients 写一个可靠的 producer + consumer,再沿 worked → partial → 开放练习三阶往上走。每个示例都把代码的某一行钉回前两章的某个概念或某个取舍——写出来的不是 API 调用,是被验证过的设计判断。
本章你将建立的 schema
- 本地用 KRaft 单 broker(无 ZooKeeper)起一个可写可读的环境,认得三类常见环境失败。
- 一个"可靠"producer 的最小配置:
acks=all+enable.idempotence=true,以及它各自挡住什么。 - 一个"不丢消息"的 consumer 循环:关 auto-commit、处理完再
commitSync(),及其投递语义后果。 - 消费侧幂等:在 at-least-once 之上用业务唯一键去重,何时它比事务 EOS 更合适。
本章代码基于 kafka-clients 4.x 的 API 编写,未在本机逐一运行。配置项名、方法签名、异常类型按官方文档与 javadoc 核对;运行命令给的是标准 quickstart 路径。把它当作"对照官方文档可直接落地"的骨架,跑通前请按你本地的 broker 地址、Java 版本核一遍。
byte[]——序列化/反序列化是必经的边界;而 commitSync()(虚线回环)发生在处理之后,这一步的位置就是后面投递语义的全部分歧所在。3.1环境准备:KRaft 单 broker + 最小依赖
一个 Maven 依赖、一个 KRaft 单 broker(无 ZooKeeper)、Java 11+ 的客户端——三样齐了就能跑本章所有示例。
客户端依赖(Maven)
只需要 kafka-clients 一个坐标,外加一个日志门面的实现(否则 SLF4J 会在启动时刷一行 "no provider" 警告)。不要引整个 kafka_2.13 服务端包——那是 broker 的依赖,客户端用不上。
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-clients</artifactId>
<version>4.0.0</version> <!-- 4.x;本章 API 在 4.0–4.3 通用 -->
</dependency>
<dependency>
<groupId>org.slf4j</groupId>
<artifactId>slf4j-simple</artifactId> <!-- 任意 SLF4J 实现即可 -->
<version>2.0.13</version>
</dependency>
kafka-clients 4.x 的客户端只要求 Java 11+(broker / Streams 需要 Java 17,但那是跑 broker 时的事,与你的应用工程无关)。Java 8 在 4.0 已被移除——若工程还卡在 Java 8,要么升级 JDK,要么把客户端停在 3.x。
本地起一个 KRaft 单 broker
4.0 起 ZooKeeper 被彻底移除,单机也走 KRaft——同一个进程既当 broker 又当 controller(process.roles=broker,controller)。最快的路径是 Docker:
# 官方镜像 apache/kafka 自带 KRaft 默认配置,开箱即用、无需 ZooKeeper
docker run -d --name kafka -p 9092:9092 apache/kafka:latest
# 验证:列出 topic(应为空列表,不报错即连通)
docker exec kafka /opt/kafka/bin/kafka-topics.sh \
--bootstrap-server localhost:9092 --list
# 建一个 4 分区、单副本的 orders topic(本机单 broker 只能 RF=1)
docker exec kafka /opt/kafka/bin/kafka-topics.sh \
--bootstrap-server localhost:9092 \
--create --topic orders --partitions 4 --replication-factor 1
不用 Docker 时,二进制发行版需要显式走 KRaft 的三步:生成集群 ID → storage format 格式化数据目录 → 启动。少了 format 这一步,broker 会因"未格式化的存储目录"直接启动失败——这是新手第一道坎。
# 1. 生成一个集群 ID(KRaft 集群的唯一标识)
KAFKA_CLUSTER_ID="$(bin/kafka-storage.sh random-uuid)"
# 2. 用该 ID 格式化数据目录(KRaft 必需,缺这步 broker 启动即失败)
# config/server.properties 已是 KRaft 单机默认(process.roles=broker,controller)
bin/kafka-storage.sh format --standalone \
--cluster-id "$KAFKA_CLUSTER_ID" \
--config config/server.properties
# 3. 启动 broker(前台运行;后台加 -daemon)
bin/kafka-server-start.sh config/server.properties
① 连接被拒 / 超时。多半是 advertised.listeners 没指向客户端能到达的地址。Docker 场景下若从宿主机连,要确保 broker advertise 的是 localhost:9092(官方镜像默认已处理);容器互联则要 advertise 容器名。客户端连不上时,第一件事是核对 broker 实际 advertise 的 host:port,而不是改客户端代码。
② 端口被占。9092 被上一个没关干净的 broker 或别的进程占用 → 启动报 Address already in use。换端口或先 docker rm -f kafka / 杀掉旧进程。
③ broker 未 format。二进制方式跳过 kafka-storage.sh format 直接 start → 报错说存储目录未格式化。KRaft 下 format 不是可选步骤。
3.2示例 1(worked):一个不丢消息的 producer + consumer
目标是一对可靠的程序:producer 保证"已确认的写入不会因单台故障丢失、且重试不produce重复",consumer 保证"处理过的记录绝不漏提交"。这两个保证不是默认值——默认配置反而会丢消息(见 04 章)。下面逐配置项落地。
可靠 producer
import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.Producer;
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;
public class ReliableProducer {
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.ACKS_CONFIG, "all"); // 等所有 ISR 确认
props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true); // 重试不produce重复、且保序
props.put(ProducerConfig.RETRIES_CONFIG, Integer.MAX_VALUE); // 配合幂等,失败就重试
// try-with-resources 保证 close() 时 flush 掉缓冲区里未发完的记录
try (Producer<String, String> producer = new KafkaProducer<>(props)) {
for (int i = 0; i < 5; i++) {
String orderId = "order-" + i;
// key = orderId:同一订单的事件落同一分区,分区内有序
ProducerRecord<String, String> record =
new ProducerRecord<>("orders", orderId, "{\"id\":\"" + orderId + "\",\"amount\":100}");
// send 异步返回 Future;回调里能拿到 broker 分配的 partition / offset
producer.send(record, (RecordMetadata md, Exception e) -> {
if (e != null) {
System.err.println("发送失败: " + e.getMessage());
} else {
System.out.printf("已确认 key=%s -> partition=%d offset=%d%n",
orderId, md.partition(), md.offset());
}
});
}
producer.flush(); // 阻塞到上面 5 条都拿到 broker 的确认
}
}
}
逐行解读(代码 → 概念,不是代码 → 语法)
min.insync.replicas(本机单 broker 演示不到,进生产必须配,见 partial 示例)。
(PID, 分区) + 序列号去重,并拒绝乱序 batch——一次解决"重试导致重复"和"重试导致乱序"两件事,对应 §2.5 投递语义。边界:它只在同一 producer 会话内去重,进程崩溃重启领新 PID 就失效。
orderId 的所有事件 murmur2 哈希到同一 partition,于是按写入顺序排列。
Future,由后台 I/O 线程批量发送。回调拿到的 partition / offset 正是 §1.2 那条日志里记录的物理坐标。flush() / close() 阻塞到缓冲区清空——漏掉它会丢掉还没发出的记录。
可靠 consumer
消费侧的两个关键决定:关掉 auto-commit、处理完一批再 commitSync()。默认的 auto-commit 按定时器提交"poll 返回的"而非"处理过的",处理到一半崩溃就丢消息(04 章头号案例)。
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.ConsumerRecords;
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 ReliableConsumer {
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-processor"); // 消费组名
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
// —— 不丢消息的两个决定 ——
props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false); // 关自动提交
props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); // 无已存 offset 时从头读
try (KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props)) {
consumer.subscribe(List.of("orders")); // 加入组,由再平衡分到若干 partition
while (true) {
// poll 既拉数据,也向 broker 证明本消费者存活
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(500));
for (ConsumerRecord<String, String> r : records) {
System.out.printf("收到 partition=%d offset=%d key=%s value=%s%n",
r.partition(), r.offset(), r.key(), r.value());
// ... 这里是真实业务:写库 / 调下游 ...
}
// 处理完整批后再同步提交 offset(at-least-once)
if (!records.isEmpty()) {
consumer.commitSync();
}
}
}
}
}
逐行解读
group.id 再跑一份,就是一个独立消费组,从头各读一遍同一份日志。
ConsumerRecords,不是单条——批的边界就是图 3.1 里那个反序列化关口。poll(Duration) 还兼任心跳证明:处理太慢、长时间不回到 poll,会被判死触发再平衡(§2.4 / 04 章再平衡风暴)。
运行 + 预期输出
两个类编译后,先跑 consumer(让它先加入组、订阅好),再跑 producer。
# 终端 A:先启动消费者(持续运行,等待消息)
mvn -q compile exec:java -Dexec.mainClass=ReliableConsumer
# 终端 B:再运行生产者,发 5 条后退出
mvn -q compile exec:java -Dexec.mainClass=ReliableProducer
终端 B(producer)按确认顺序打印——注意不同 key 被路由到不同 partition:
已确认 key=order-0 -> partition=2 offset=0
已确认 key=order-1 -> partition=0 offset=0
已确认 key=order-2 -> partition=2 offset=1
已确认 key=order-3 -> partition=1 offset=0
已确认 key=order-4 -> partition=3 offset=0
# 具体 partition 取决于 murmur2(key) % 4,因 key 而定、可复现;offset 在各分区内从 0 起。
终端 A(consumer)收到这 5 条。同一 partition 内按 offset 有序,跨 partition 之间无全局顺序:
收到 partition=0 offset=0 key=order-1 value={"id":"order-1","amount":100}
收到 partition=1 offset=0 key=order-3 value={"id":"order-3","amount":100}
收到 partition=2 offset=0 key=order-0 value={"id":"order-0","amount":100}
收到 partition=2 offset=1 key=order-2 value={"id":"order-2","amount":100}
收到 partition=3 offset=0 key=order-4 value={"id":"order-4","amount":100}
# partition 间到达顺序不保证;partition=2 内 offset 0 必在 1 之前——这就是"顺序只在分区内"。
consumer 跑完、不重启它,把 producer再运行一遍(又发 5 条相同 key)。这次 consumer 会收到这 5 条吗?如果先 Ctrl-C 停掉 consumer、重跑 producer、再启动 consumer,结果一样吗?
展开答案(先停 10 秒再点)
第一种(consumer 不停):会收到。新发的 5 条追加到各分区日志尾部,consumer 的游标继续往后 poll 到它们。这 5 条的 offset 接着上次往后排(比如 partition=2 收到 offset 2、3)。
第二种(停掉再重启):也只收到这新的 5 条,不会重收旧的。因为上一轮 commitSync() 已经把 offset 提交到 __consumer_offsets,同 group.id 的 consumer 重启后从已提交位置续读,不是从头。auto.offset.reset=earliest 只在"该组从来没提交过 offset"时才从头——一旦提交过,它不生效。
这道题指向:offset 归消费组持有、跨重启存活(§1.2),以及"读完不删"——旧的 5 条始终在日志里,换个新 group.id 仍能从头再读一遍。
3.3示例 2(partial):补上生产级的三个决策点
worked 版在本机单 broker 上能跑,但有三处是本机演示不到、进生产必须做的决策。下面的框架把这三处留成 TODO——它们不是填空,是要你在两个真实备选间做判断并说出代价。先自己定,再展开对照决策路径。
// —— Producer 端 ——
props.put(ProducerConfig.ACKS_CONFIG, "all");
props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true);
// TODO ①:broker 端 topic 该配 min.insync.replicas 为几?(RF=3 的前提下)
// 提示:acks=all 的强度被当前 ISR 大小绑架,这个值决定 ISR 缩到多少时拒写。
// —— 这是 broker/topic 配置,不是 producer 属性,用 kafka-configs.sh 设。
// —— Consumer 端 ——
props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(500));
for (ConsumerRecord<String, String> r : records) {
handle(r.value());
// TODO ②:commit 放这里(每条一提交)还是放 for 循环之后(每批一提交)?
// 两种都能跑通,但投递语义和吞吐不同。选一个,并说出代价。
}
// TODO ②(另一处):还是放这里?
}
// TODO ③:orders 要求"同一订单的事件严格按 创建→支付→发货 顺序被消费"。
// 现在靠 key=orderId 的默认分区器够不够?要不要自定义 Partitioner?
TODO ① 参考答案 + 决策路径
设 min.insync.replicas=2(RF=3 时)。决策路径:acks=all 只保证"等当前 ISR 全部确认",但 ISR 会动态收缩——两个 follower 掉队后 ISR={leader},"all" 退化成一台,等价 acks=1 且 producer 毫无察觉(§2.3 的头号陷阱)。设成 2 才能在 ISR 缩到 1 时让写入被拒(NotEnoughReplicas)而非悄悄降级。
代价:RF=3 + min.insync.replicas=2 下,挂 1 台仍可写、挂 2 台该分区停写——这是故意用可用性换持久性。设成 3 则任一台故障就停写(太脆),设成 1 则等于没设(退化成上面的陷阱)。规则:min.insync.replicas = RF - 1 是"容忍一台故障还能写"的平衡点。
TODO ② 参考答案 + 决策路径
放 for 循环之后(每批一提交)是默认选择。两种都是 at-least-once(都先处理后提交),区别在提交频率:
- 每批一提交:一次
commitSync()覆盖整批,提交开销小、吞吐高;代价是崩溃时整批重放(重复窗口 = 一批)。 - 每条一提交:重复窗口缩到一条,但每条都同步往 broker 发一次提交、阻塞等确认,吞吐大幅下降。
决策路径:既然消费侧已经要做幂等(见 开放练习),重放几条不产生副作用,就没必要为缩小重复窗口牺牲吞吐——每批一提交 + 幂等是标准组合。只有当处理极慢、单条都很贵时才考虑更细的提交粒度(或用 commitAsync() 配合定期 commitSync())。唯一不能选的是把 commit 挪到 handle() 之前——那会从重复翻成丢失。
TODO ③ 参考答案 + 决策路径
不需要自定义 Partitioner。key=orderId + 默认分区器已经保证"同一 orderId → murmur2 哈希 → 同一 partition → 分区内按写入顺序"(§2.2)。"创建→支付→发货"只要按这个先后 send,且 producer 开了 enable.idempotence=true(幂等顺带拒绝乱序 batch),到达消费端就是这个顺序。
什么时候才需要自定义 Partitioner:默认按 key 哈希均摊,但某些场景要把"一组相关 key"压进同一分区(如同一商户的所有订单要单分区串行处理),或要规避热点 key 倾斜——这时才写 Partitioner。代价提醒:无论默认还是自定义,partition 数一旦上线就不能减,给有 key 的 topic 加分区会打乱已有 key→分区映射、破坏历史顺序(§2.2)。所以顺序设计的重点在"key 选得对 + 分区数预估够",而非分区器本身。
3.4示例 3(开放练习):幂等的订单事件消费者
实现一个消费者:从 orders topic 消费订单事件,把每个订单写入数据库,保证 Kafka 的重复投递不产生重复订单。验收:人为让某一批在提交 offset 前崩溃、重启后重放同一批,数据库里那些订单的记录数不变。
必须用到 01 章的 offset 由消费者维护、consumer group、record 的 key 三个概念,以及 02 章的一个权衡:at-least-once + 消费侧幂等 vs 事务 EOS——你要选其一并说清为什么。
这是真实工作里最常见的形态。worked 示例已经给了"不丢"(at-least-once),但 at-least-once 的代价是会重复:再平衡、崩溃重启都会重放未提交的那批。"不产生重复订单"这件事,Kafka 的事务 EOS 解决不了——因为写数据库是外部副作用,不在 Kafka 事务范围内(§2.5 EOS 的真实边界)。所以正确的下手处是消费侧幂等:用订单的业务唯一键让"重复处理"变成"无害的重复写"。
脚手架递减:你现在站在第三阶
参考实现 + 关键决策说明(写完自己版本再展开)
核心思路:把 worked 版的 consumer 改成"在同一个数据库事务里,用订单 ID 做幂等键 upsert/INSERT-IGNORE"。Kafka 这侧仍是 at-least-once(处理后提交 offset),重复由数据库的唯一约束吸收。
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.common.serialization.StringDeserializer;
import javax.sql.DataSource;
import java.sql.Connection;
import java.sql.PreparedStatement;
import java.time.Duration;
import java.util.List;
import java.util.Properties;
public class IdempotentOrderConsumer {
private final DataSource ds; // 你的连接池
public IdempotentOrderConsumer(DataSource ds) { this.ds = ds; }
public void run() {
Properties props = new Properties();
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(ConsumerConfig.GROUP_ID_CONFIG, "order-db-writer");
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false); // 手动提交,处理后才提交
try (KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props)) {
consumer.subscribe(List.of("orders"));
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(500));
if (records.isEmpty()) continue;
// 整批在一个 DB 事务里写;幂等键 = 订单 ID(来自 record.key())
try (Connection conn = ds.getConnection()) {
conn.setAutoCommit(false);
// ON CONFLICT DO NOTHING:同一 order_id 重复投递时第二次写入被静默忽略
String sql = "INSERT INTO orders(order_id, payload) VALUES (?, ?) "
+ "ON CONFLICT (order_id) DO NOTHING";
try (PreparedStatement ps = conn.prepareStatement(sql)) {
for (ConsumerRecord<String, String> r : records) {
ps.setString(1, r.key()); // order_id 作幂等键
ps.setString(2, r.value());
ps.addBatch();
}
ps.executeBatch();
}
conn.commit(); // DB 先落库
}
// DB 提交成功后才提交 Kafka offset。崩在两者之间 → 重放整批,
// 但 DB 的 ON CONFLICT DO NOTHING 吸收重复,记录数不变。
consumer.commitSync();
}
}
}
}
关键决策:为什么用消费侧幂等而非事务 EOS。
- 事务 EOS 管不到外部副作用。Kafka 事务(
transactional.id+sendOffsetsToTransaction)只能把"写回 Kafka 的输出 + offset 提交"做成原子——本练习的副作用是写数据库,发生在 Kafka 之外,事务回滚撤不掉已发生的 DB 写(§2.5 的核心边界)。 - 幂等键把"精确一次"下沉到有副作用的地方。唯一能让"写库恰好一次"成立的位置,是写库本身——靠
order_id的唯一约束 +ON CONFLICT DO NOTHING。这样 Kafka 侧保持简单的 at-least-once,正确性由 DB 兜底。 - 顺序:DB 提交在前,offset 提交在后。这个先后不能反——先提交 offset 再写库,崩在中间就丢了那批(at-most-once)。先库后 offset,崩在中间只会重放,被幂等吸收。
什么时候反而该用事务 EOS:当副作用就是写回另一个 Kafka topic("消费 A → 转换 → 生产 B"的纯 Kafka 闭环),事务能把 B 的输出和 A 的 offset 原子提交,比"消费侧幂等 + 自己维护去重表"更干净。判别点就一句:副作用落在 Kafka 内还是 Kafka 外。
现状延伸:share group(KIP-932,4.2+)让 Kafka 有了队列语义
本章的消费模型自始至终是"一个 partition 只归组内一个 consumer"——并行度被 分区数卡死,且天然是 at-least-once。截至 2026-06,Kafka 4.2(2026-02)GA 的 share group("Queues for Kafka")在日志之上加了一层队列语义:同一 topic 的记录可被同组多个 consumer 按记录分发并逐条 ack,不再受"消费者数 ≤ 分区数"限制——更接近传统 MQ 的工作队列。
share group 用的是 KafkaShareConsumer,poll() 后对每条记录调 acknowledge(record, AcknowledgeType.ACCEPT / RELEASE / REJECT),再 commitSync() 提交这些 ack——ack 状态由 broker 端维护,而不是单一 offset 游标。它仍架在同一条分区日志之上(数据没删、可被普通消费组重放),只是把"谁处理哪条、哪条处理失败要重投"的记账搬到了 broker。需要"工作队列 + 逐条重试 + 并行度不绑分区数"时它是原生答案;需要"可重放 + 严格分区顺序"时仍用经典 consumer group。两者共存,按语义选,不是替代关系。
§本章 self-check(动手层面)
先合上教程,把你能想到的答案写在编辑器里。写完再点开对照——直接点开等于把这一节当再读一遍。
- worked 版 producer 配了
acks=all,为什么说它"单独不够"?要在哪一侧、补哪个配置才补得上?(指出配置归 producer 还是 broker/topic) - 把可靠 consumer 里的
commitSync()从 for 循环之后挪到handle()之前,投递语义从什么变成什么?什么时候会观察到差异? - 开放练习里"DB 提交"和"Kafka offset 提交"两步,为什么顺序必须是先 DB 后 offset?反过来会怎样?
- (设计题)同样要"消费后写下游且不重复",什么情况下你选事务 EOS、什么情况下选消费侧幂等?判别的那一句话是什么?
答案(先做完再展开)
acks=all只保证"等当前 ISR 全部确认",而 ISR 会动态收缩——缩到只剩 leader 时 "all" 就是一台,悄悄退化成acks=1且 producer 无感知。补在 broker/topic 侧的min.insync.replicas=2(RF=3 时),让 ISR 缩到 1 时写入被拒(NotEnoughReplicas)而非降级。这是 topic 配置,不是 producer 属性,用kafka-configs.sh设。- 从 at-least-once 变成 at-most-once。差异在崩溃/再平衡时出现:先提交后处理时,提交完那批还没处理完就崩溃,重启从已提交位置之后续读——那批被跳过(丢失)。先处理后提交则只会重放(重复)。"感觉更安全"的早提交才是危险的。
- 先 DB 后 offset:崩在两步之间只会让整批被重放,DB 的唯一约束(
ON CONFLICT DO NOTHING)吸收重复,记录数不变。反过来先提交 offset 再写库:崩在中间时 offset 已前移、那批的 DB 写没发生,重启不会再读到它们——丢数据(at-most-once)。 - 判别那一句:副作用落在 Kafka 内还是 Kafka 外。落在 Kafka 内(消费 A → 生产 B 的纯闭环)→ 事务 EOS,把 B 输出和 A 的 offset 原子提交。落在 Kafka 外(写库、调支付、发邮件)→ 事务管不到,必须在副作用那一侧用业务幂等键去重,Kafka 侧保持 at-least-once。
把 worked 示例改成事务 EOS 版(纯 Kafka 闭环)
把示例 1 改成"消费 orders → 计算 → 生产到 orders-enriched"的纯 Kafka管道,并用事务让"输出写入 + 消费 offset 提交"原子。producer 要配 transactional.id 并调 initTransactions() / beginTransaction() / sendOffsetsToTransaction(...) / commitTransaction();consumer 要把 offset 提交交给事务(而不是自己 commitSync()),并设 isolation.level=read_committed。想清楚:为什么这里能用 EOS,而开放练习的写库场景不能?
提示(卡住再展开)
关键差别在副作用的位置——这个挑战的输出是另一个 Kafka topic(Kafka 内),所以事务覆盖得到;开放练习的输出是数据库(Kafka 外),事务覆盖不到。事务里 offset 不能再用 consumer.commitSync() 提交,而要走 producer.sendOffsetsToTransaction(offsets, consumer.groupMetadata()),让 offset 提交和输出写入进同一个事务、一起 commit 或一起 abort。read_committed 的消费者会跳过 aborted 数据、且不越过 LSO(§2.5)。02 章那段事务 producer 代码可以直接当骨架。