Kafka学习笔记(二):Docker与Java入门实战
上一篇把 Kafka 理解成了“多个读者共享的一本事件账本”。这一篇直接动手:启动 Kafka,创建订单 Topic,写入事件,再让 Java 消费者读取它。
这次实验的目标是:能看到消息进入哪个分区,拿到什么 Offset,以及业务处理后在哪里提交进度。
系列导航:
-
核心概念与架构
-
Docker 与 Java 入门实战
-
消费组、Offset 与 Rebalance
-
消息可靠性、幂等与事务
-
面试高频问题总结
一、环境准备
| 工具 |
本文的实验约定 |
| 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) { } } }
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 为什么会造成重复消费或漏处理。