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 版本核一遍。

动作 类型 ProducerRecord 业务对象 <K,V> 序列化 broker 分区日志 append 到尾部 byte[] poll ConsumerRecords 一批记录 <K,V> 遍历 业务处理 handle(value) V 处理后 commitSync() 提交 offset 回 broker
图 3.1一条记录从业务对象到落盘、再回到业务对象的完整类型链。 注意:两端(朱红)是 JVM 里的对象,broker 上只存 byte[]——序列化/反序列化是必经的边界;而 commitSync()(虚线回环)发生在处理之后,这一步的位置就是后面投递语义的全部分歧所在。

3.1环境准备:KRaft 单 broker + 最小依赖

一个 Maven 依赖、一个 KRaft 单 broker(无 ZooKeeper)、Java 11+ 的客户端——三样齐了就能跑本章所有示例。

客户端依赖(Maven)

只需要 kafka-clients 一个坐标,外加一个日志门面的实现(否则 SLF4J 会在启动时刷一行 "no provider" 警告)。不要引整个 kafka_2.13 服务端包——那是 broker 的依赖,客户端用不上。

pom.xml(片段) properties
<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>
版本 / Java 要求

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:

起 broker(Docker,推荐) bash
# 官方镜像 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 会因"未格式化的存储目录"直接启动失败——这是新手第一道坎。

起 broker(二进制 + KRaft,无 ZooKeeper) bash
# 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

ReliableProducer.java Java
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 的确认
        }
    }
}

逐行解读(代码 → 概念,不是代码 → 语法)

acks=all把"写成功"的定义钉在"当前 ISR 全部确认"上,对应 §2.3 副本与 ISR。但它单独不够——ISR 缩到只剩 leader 时 "all" 就是一台,持久性真正的闸门是 broker 端的 min.insync.replicas(本机单 broker 演示不到,进生产必须配,见 partial 示例)。
enable.idempotence开启后 broker 用 (PID, 分区) + 序列号去重,并拒绝乱序 batch——一次解决"重试导致重复"和"重试导致乱序"两件事,对应 §2.5 投递语义。边界:它只在同一 producer 会话内去重,进程崩溃重启领新 PID 就失效。
key=orderId是把"需要顺序的记录钉在同一分区"的落点,对应 §1.3 顺序只在分区内——同一 orderId 的所有事件 murmur2 哈希到同一 partition,于是按写入顺序排列。
send是异步的:它把记录放进缓冲区就返回 Future,由后台 I/O 线程批量发送。回调拿到的 partition / offset 正是 §1.2 那条日志里记录的物理坐标。flush() / close() 阻塞到缓冲区清空——漏掉它会丢掉还没发出的记录。

可靠 consumer

消费侧的两个关键决定:关掉 auto-commit、处理完一批再 commitSync()。默认的 auto-commit 按定时器提交"poll 返回的"而非"处理过的",处理到一半崩溃就丢消息(04 章头号案例)。

ReliableConsumer.java Java
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把这个消费者归进一个消费组,组内每个 partition 只分给一个成员——对应 §1.4 并行度 = 分区数。换一个 group.id 再跑一份,就是一个独立消费组,从头各读一遍同一份日志。
poll返回的是一批 ConsumerRecords,不是单条——批的边界就是图 3.1 里那个反序列化关口。poll(Duration) 还兼任心跳证明:处理太慢、长时间不回到 poll,会被判死触发再平衡(§2.4 / 04 章再平衡风暴)。
commitSync放在 for 循环之后,决定了投递语义是at-least-once:先处理后提交,崩溃只会让上一批被重放(重复),不会丢。把它挪到处理之前就翻成 at-most-once(丢失)——这正是 §2.5 那条"先提交后处理把失败模式从重复翻成丢失"。

运行 + 预期输出

两个类编译后,先跑 consumer(让它先加入组、订阅好),再跑 producer。

运行命令 bash
# 终端 A:先启动消费者(持续运行,等待消息)
mvn -q compile exec:java -Dexec.mainClass=ReliableConsumer

# 终端 B:再运行生产者,发 5 条后退出
mvn -q compile exec:java -Dexec.mainClass=ReliableProducer

终端 B(producer)按确认顺序打印——注意不同 key 被路由到不同 partition:

producer 预期输出 bash
已确认 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 之间无全局顺序:

consumer 预期输出 bash
收到 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——它们不是填空,是要你在两个真实备选间做判断并说出代价。先自己定,再展开对照决策路径。

ProductionConfig.java(待补) Java
// —— 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 的真实边界)。所以正确的下手处是消费侧幂等:用订单的业务唯一键让"重复处理"变成"无害的重复写"。

脚手架递减:你现在站在第三阶

示例 1 · worked 完整代码全给 逐行解读 示例 2 · partial 框架给,留决策点 3 个 TODO 示例 3 · open 自己写,参考实现折叠 作者承担:多 作者承担:中 读者承担:多 脚手架(作者给的代码量)递减 →
图 3.2三阶脚手架:实心框的高度近似"作者直接给出的代码量",逐阶变矮。 注意:递减的是脚手架,递增的是你的承担——到第三阶(朱红)参考实现是折叠的,先自己写完再展开,对照才有意义;直接看答案等于把练习做成了阅读。
参考实现 + 关键决策说明(写完自己版本再展开)

核心思路:把 worked 版的 consumer 改成"在同一个数据库事务里,用订单 ID 做幂等键 upsert/INSERT-IGNORE"。Kafka 这侧仍是 at-least-once(处理后提交 offset),重复由数据库的唯一约束吸收。

IdempotentOrderConsumer.java(参考) Java
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 没有推翻"日志"模型

share group 用的是 KafkaShareConsumer,poll() 后对每条记录调 acknowledge(record, AcknowledgeType.ACCEPT / RELEASE / REJECT),再 commitSync() 提交这些 ack——ack 状态由 broker 端维护,而不是单一 offset 游标。它仍架在同一条分区日志之上(数据没删、可被普通消费组重放),只是把"谁处理哪条、哪条处理失败要重投"的记账搬到了 broker。需要"工作队列 + 逐条重试 + 并行度不绑分区数"时它是原生答案;需要"可重放 + 严格分区顺序"时仍用经典 consumer group。两者共存,按语义选,不是替代关系。

§本章 self-check(动手层面)

先合上教程,把你能想到的答案写在编辑器里。写完再点开对照——直接点开等于把这一节当再读一遍。

  1. worked 版 producer 配了 acks=all,为什么说它"单独不够"?要在哪一侧、补哪个配置才补得上?(指出配置归 producer 还是 broker/topic)
  2. 把可靠 consumer 里的 commitSync() 从 for 循环之后挪到 handle() 之前,投递语义从什么变成什么?什么时候会观察到差异?
  3. 开放练习里"DB 提交"和"Kafka offset 提交"两步,为什么顺序必须是先 DB 后 offset?反过来会怎样?
  4. (设计题)同样要"消费后写下游且不重复",什么情况下你选事务 EOS、什么情况下选消费侧幂等?判别的那一句话是什么?
答案(先做完再展开)
  1. 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 设。
  2. 从 at-least-once 变成 at-most-once。差异在崩溃/再平衡时出现:先提交后处理时,提交完那批还没处理完就崩溃,重启从已提交位置之后续读——那批被跳过(丢失)。先处理后提交则只会重放(重复)。"感觉更安全"的早提交才是危险的。
  3. 先 DB 后 offset:崩在两步之间只会让整批被重放,DB 的唯一约束(ON CONFLICT DO NOTHING)吸收重复,记录数不变。反过来先提交 offset 再写库:崩在中间时 offset 已前移、那批的 DB 写没发生,重启不会再读到它们——丢数据(at-most-once)。
  4. 判别那一句:副作用落在 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 代码可以直接当骨架。