Kafka学习笔记(二):Docker与Java入门实战

上一篇把 Kafka 理解成了“多个读者共享的一本事件账本”。这一篇直接动手:启动 Kafka,创建订单 Topic,写入事件,再让 Java 消费者读取它。

这次实验的目标是:能看到消息进入哪个分区,拿到什么 Offset,以及业务处理后在哪里提交进度。

系列导航:

  1. 核心概念与架构
  2. Docker 与 Java 入门实战
  3. 消费组、Offset 与 Rebalance
  4. 消息可靠性、幂等与事务
  5. 面试高频问题总结

一、环境准备

工具 本文的实验约定
Docker 可运行 Linux 容器;Windows 可使用已启动的 Docker Desktop
Kafka 镜像 apache/kafka:4.0.0
Java JDK 17
Maven 3.6.3 或更高版本,用于编译与启动 Java 示例
客户端位置 Java 程序在宿主机运行,Kafka 在本机容器运行

固定 4.0.0 是为了让系列示例一致,不表示它是最新版本。Kafka 4.0 使用 KRaft,本文不需要额外启动 ZooKeeper。官方升级说明与官方 Docker 指南可用于核对环境要求。

下面的 Docker 命令都写成单行,在 PowerShell 或 Bash 中都可以执行,避免混用续行符。

1
2
3
docker version
java -version
mvn -version

确认 docker version 能返回 Server 信息,并且 mvn -version 使用的 Java 也是 JDK 17。只看 java -version,可能忽略 Maven 的 JAVA_HOME 仍指向旧 JDK。本文使用的 exec 插件要求 Maven 至少为 3.6.3,见插件版本要求。


二、启动一个 Kafka 节点

2.1 使用官方镜像

1
docker run -d --name kafka-learning -p 127.0.0.1:9092:9092 apache/kafka:4.0.0

这个实验使用官方镜像的默认单节点配置,端口只映射到本机。镜像用法来自官方 Docker 指南,4.0.0 标签可在官方版本列表中查到。

检查日志和服务是否就绪:

1
2
docker logs --tail 100 kafka-learning
docker exec kafka-learning /opt/kafka/bin/kafka-topics.sh --bootstrap-server localhost:9092 --list

首次启动可能需要等待一会儿。容器状态显示 Up,不代表 Kafka 已经可以处理请求;以命令行成功连接为准。没有 Topic 时,列表为空是正常结果。

2.2 两个 localhost,各指什么?

1
2
3
4
5
宿主机上的 Java
localhost:9092 → 本机端口映射 → Kafka 容器

docker exec 中的 Kafka CLI
localhost:9092 → 容器内部的 Kafka

本文只覆盖这两条访问路径。如果把 Java 应用也放进另一个容器,应用里的 localhost 会指向应用容器本身,需要重新设计监听器和对外广播地址。

bootstrap.servers 是发现集群的入口,客户端随后会连接元数据中广播的 Broker 地址。 所以排查连接问题时,还要检查 advertised.listeners 是否能从客户端所在位置访问。客户端配置文档解释了这种发现过程。


三、命令行完成第一轮生产与消费

3.1 创建订单 Topic

1
docker exec kafka-learning /opt/kafka/bin/kafka-topics.sh --bootstrap-server localhost:9092 --create --if-not-exists --topic order-events --partitions 3 --replication-factor 1

这里创建 3 个分区、每个分区 1 个副本。如果 Topic 已经存在,--if-not-exists 不会把它改成三个分区,因此要查看实际配置:

1
docker exec kafka-learning /opt/kafka/bin/kafka-topics.sh --bootstrap-server localhost:9092 --describe --topic order-events

关注三个字段:

  • PartitionCount:分区数量。
  • Leader:负责该分区的 Broker ID。
  • Replicas、Isr:副本列表与同步副本列表。

这里只有一个 Broker,副本数设成 3 会因节点不足而失败。三个分区可以练习并行消费,但单副本无法提供节点故障时的数据冗余。

Topic 管理命令的基本用法可参考官方快速入门。

3.2 启动带 Key 的生产者

1
docker exec -it kafka-learning /opt/kafka/bin/kafka-console-producer.sh --bootstrap-server localhost:9092 --topic order-events --property parse.key=true --property key.separator=:

进入交互界面后,逐行输入:

1
2
3
O1001:{"eventId":"E1001","orderId":"O1001","type":"CREATED"}
O1001:{"eventId":"E1002","orderId":"O1001","type":"PAID"}
O1002:{"eventId":"E1003","orderId":"O1002","type":"CREATED"}

第一个冒号前的内容是 Key,后面的内容是 Value。O1001 的两次事件使用同一个 Key;不要预设它一定进入 Partition 0,观察实际结果即可。

3.3 在另一个终端启动消费者

1
docker exec -it kafka-learning /opt/kafka/bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic order-events --group order-cli-demo-v1 --from-beginning --property print.key=true --property print.partition=true --property print.offset=true

可以看到每条消息的 Key、Partition、Offset 和 Value。不同分区的显示先后顺序可能交错,这是正常的。

如果同一个组已有有效的已提交进度,--from-beginning 不会覆盖这个进度。 想独立观察保留中的历史记录,可以换一个全新的组名,比如 order-cli-demo-v2。具体进度规则会在下一篇展开。


四、创建一个最小 Java 项目

新建目录 kafka-learning-demo,结构如下:

1
2
3
4
5
kafka-learning-demo/
├── pom.xml
└── src/main/java/demo/
├── OrderProducer.java
└── OrderConsumer.java

pom.xml 内容:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 https://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<groupId>demo</groupId>
<artifactId>kafka-learning-demo</artifactId>
<version>1.0-SNAPSHOT</version>
<properties>
<maven.compiler.release>17</maven.compiler.release>
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
</properties>
<dependencies>
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-clients</artifactId>
<version>4.0.0</version>
</dependency>
</dependencies>
<build>
<plugins>
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-compiler-plugin</artifactId>
<version>3.13.0</version>
</plugin>
<plugin>
<groupId>org.codehaus.mojo</groupId>
<artifactId>exec-maven-plugin</artifactId>
<version>3.5.0</version>
</plugin>
</plugins>
</build>
</project>

这次直接使用原生客户端,方便观察 send()、poll() 和 commitSync() 的职责。Value 使用 JSON 字符串,暂时不引入对象序列化框架。


五、生产者:发送成功要看结果

保存为 src/main/java/demo/OrderProducer.java:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
package demo;

import java.util.Properties;
import java.util.UUID;
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;

public class OrderProducer {
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());
props.put(ProducerConfig.ACKS_CONFIG, "all");
props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, "true");

try (KafkaProducer<String, String> producer = new KafkaProducer<>(props)) {
String eventId = UUID.randomUUID().toString();
String value = String.format(
"{\"eventId\":\"%s\",\"orderId\":\"O1001\",\"type\":\"CREATED\"}",
eventId);
ProducerRecord<String, String> record =
new ProducerRecord<>("order-events", "O1001", value);

RecordMetadata result = producer.send(record).get();
System.out.printf("发送成功:topic=%s, partition=%d, offset=%d%n",
result.topic(), result.partition(), result.offset());
}
}
}

5.1 为什么这里调用 get()?

send() 通常先把记录放入客户端缓冲区,实际发送由后台线程完成。调用 .get() 等待这次发送的最终结果,便于入门时定位失败。KafkaProducer API说明了这种异步行为。

逐条等待会降低批量发送的效率。实际业务可以使用回调观察成功或失败,但不能只调用 send() 就认定消息已经进入 Broker。

5.2 每次启动为什么生成新 eventId?

这是为了让重复运行示例代表一个新的演示事件。真实业务中,同一次事件的应用层重试必须复用原来的 eventId,否则消费者无法通过这个 ID 去重。

Kafka 生产者幂等也不能把两次业务上主动调用 send() 自动合并。第四篇会解释它的作用边界。


六、消费者:先处理,再提交

保存为 src/main/java/demo/OrderConsumer.java:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
package demo;

import java.time.Duration;
import java.util.List;
import java.util.Properties;
import java.util.concurrent.atomic.AtomicBoolean;
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.errors.WakeupException;
import org.apache.kafka.common.serialization.StringDeserializer;

public class OrderConsumer {
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-java-demo-v1");
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");
props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, "100");
props.put(ConsumerConfig.GROUP_PROTOCOL_CONFIG, "classic");

AtomicBoolean stopping = new AtomicBoolean(false);
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
Thread shutdownHook = new Thread(() -> {
stopping.set(true);
consumer.wakeup();
});
Runtime.getRuntime().addShutdownHook(shutdownHook);

try (consumer) {
consumer.subscribe(List.of("order-events"));
while (!stopping.get()) {
ConsumerRecords<String, String> records =
consumer.poll(Duration.ofSeconds(1));
for (ConsumerRecord<String, String> record : records) {
process(record);
}
if (!records.isEmpty()) {
consumer.commitSync();
}
}
} catch (WakeupException e) {
if (!stopping.get()) {
throw e;
}
} finally {
try {
Runtime.getRuntime().removeShutdownHook(shutdownHook);
} catch (IllegalStateException ignored) {
// JVM 已经进入关闭流程。
}
}
}

private static void process(ConsumerRecord<String, String> record) {
// 演示处理;真实业务应在事务落库成功后才返回。
System.out.printf("处理订单:partition=%d, offset=%d, key=%s, value=%s%n",
record.partition(), record.offset(), record.key(), record.value());
}
}

这段代码遵守一个约定:本轮 poll() 返回的记录全部同步处理成功,才提交本轮进度。 process() 或提交抛出异常时,程序退出,不会捕获错误后继续提交更靠后的进度。已执行但未提交的记录,重启后仍可能重复。

通过 wakeup() 通知消费线程退出,是官方客户端文档提供的关闭方式。KafkaConsumer 不能供多个线程同时操作,wakeup() 是用于跨线程通知的例外。

示例中的打印不代表持久化业务效果,也没有实现 JSON 校验、数据库事务和去重。下一步接入业务时,应同时补上这些能力。


七、运行与观察

在项目根目录执行:

1
2
mvn compile
mvn exec:java -Dexec.mainClass=demo.OrderConsumer

在另一个终端进入同一个目录,执行:

1
mvn exec:java -Dexec.mainClass=demo.OrderProducer

第一次启动新的 Java 消费组时,会读取仍保留的命令行事件,然后读取 Java 生产者新写入的事件。生产者输出发送结果,消费者输出对应的 Key、分区和 Offset。

查看消费进度:

1
docker exec kafka-learning /opt/kafka/bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 --describe --group order-java-demo-v1
字段 如何理解
CURRENT-OFFSET 这个组已经提交的下一读取位置
LOG-END-OFFSET 日志末尾的下一位置
LAG 两者之间的 Offset 差值

对于本次无事务、无压缩清理的简单日志,LAG 可以近似理解成尚未提交进度覆盖的记录数量。更复杂的日志中,Offset 可能有间隙,不能一律把这个值当作待处理业务事件的精确条数。


八、常见问题排查

现象 先检查什么
Docker 找不到 Server Docker Desktop 是否启动,是否使用 Linux 容器
Kafka 启动但 CLI 连不上 查看日志,等待就绪,再检查端口是否占用
Java 报类版本或编译错误 Java 与 Maven 使用的 JDK 是否都是 17
earliest 没重放旧消息 当前组是否已经提交过进度
启动多个消费者,有的没输出 分区分配、Key 是否集中到某个分区
设置 acks=all 仍只有一份数据 Topic 的副本数仍是 1;ACK 配置不会创建副本

实验完成后可以停掉容器,之后再启动继续学习:

1
2
docker stop kafka-learning
docker start kafka-learning

本文没有挂载持久化卷,删除并重建容器会丢失这次实验的容器内数据。后续需要长期保存数据时,再配置独立的数据卷。

下一篇使用这套环境观察消费组如何分工,以及 Offset 为什么会造成重复消费或漏处理。