卡夫卡(Kafka)从入门到实践:超详细学习指南
读完这篇博客,你将能够:用 Docker 一键部署 Kafka 集群,用 Spring Boot 收发消息,甚至传输文件。
源码链接:Misty/Kafka
一、Kafka 是什么?为什么需要它?
想象一个场景:你在外卖 App 上点了一份麻辣烫。
你下单 → 商家接单 → 骑手取餐 → 配送中 → 你收到餐
这些步骤环环相扣,每个环节都在等上一个完成。如果某个环节卡住了,整个流程就断了——这就是同步调用的痛点。
现在换个思路:中间加一个消息中转站,每个环节只跟中转站打交道。
你下单 ──→ 消息中转站 ──→ 商家接单 骑手GPS ──→ ──→ 配送追踪 商家出餐 ──→ ──→ 用户通知
-
每个环节独立工作,互不阻塞
-
新加一个"优惠券系统"?直接从中转站订阅就行,不用改任何现有代码
-
某个环节挂了?消息留在中转站里,恢复了继续处理
这个"消息中转站",就是 Kafka。
Kafka 最初由 LinkedIn 开发,现在由 Apache 基金会维护,是全球最流行的分布式消息系统。LinkedIn 每天用它处理超过 7 万亿条消息。
二、核心概念(用人话解释)
2.1 几个关键角色
| 概念 | 一句话解释 | 类比 |
|---|---|---|
| Producer(生产者) | 发消息的人 | 外卖 App 把订单扔进系统 |
| Consumer(消费者) | 收消息的人 | 商家从系统拉取订单 |
| Broker(节点) | Kafka 服务实例 | 中转站的一个分拣窗口 |
| Topic(主题) | 消息的类别 | 不同的传送带,订单一条、支付一条、通知一条 |
| Partition(分区) | Topic 的物理分片 | 一条传送带拆成 3 段,并行运转 |
| Offset(偏移量) | 消息的唯一编号 | 快递单号,消费者凭它记住"读到哪了" |
| Consumer Group(消费者组) | 一群消费者共享进度 | 三个店员抢一传送带的活儿,自动分工 |
2.2 一张图看懂消息流转
┌──────────┐ ┌──────────┐
│ Producer │ │ Consumer │
│ │ send("订单1") │ │
└────┬─────┘ └────┬─────┘
│ │
│ 消息写入 监听拉取新消息 │
▼ ▼
┌──────────────────────────────────────────────────────┐
│ Kafka 集群 │
│ │
│ Topic: order-topic │
│ ┌──────────┐ ┌──────────┐ ┌──────────┐ │
│ │Partition0│ │Partition1│ │Partition2│ │
│ │ offset 0 │ │ offset 1 │ │ offset 2 │ │
│ │ offset 3 │ │ offset 4 │ │ offset 5 │ │
│ └──────────┘ └──────────┘ └──────────┘ │
│ │
│ Broker1 Broker2 Broker3 Broker4 │
└──────────────────────────────────────────────────────┘
2.3 两个最重要的细节
分区策略——消息发给谁?
-
不指定 key:轮询,消息均匀分布到各个分区
-
指定 key:对 key 取 hash,相同 key 永远进同一分区,保证消费顺序
举个栗子:你把 userId 作为 key,那么同一个用户的所有操作消息会严格按时间顺序被消费。
Consumer Group——多人怎么分活?
-
一个分区只能被组内一个消费者消费(避免两人抢同一个活)
-
消费者数量 > 分区数量?多出来的人闲着
-
消费者数量 ≤ 分区数量?自动平均分配,有人挂了自动重新分配
三、环境搭建:Docker 一键部署 4 节点集群
3.1 ZooKeeper 与 KRaft 完整对比(通俗版)
一个故事讲清楚
Kafka 集群就像一个大型物流中转站,多个 Broker 是各个分拣窗口。窗口之间需要知道:谁是总调度?哪个窗口负责哪些包裹?有没有窗口宕机了?
ZooKeeper 模式(Kafka 2.x 及以前)
Kafka 请了一个叫 ZooKeeper 的外部管理员来管理这些事。管理员自己就是一个 3 人小团队(ZooKeeper 集群),负责:
-
记录谁是总调度窗口(Controller 选举)
-
记录每个包裹被分到哪个窗口(分区分配信息)
-
监控哪个窗口还活着(心跳检测)
┌─────────────────────────────────────────────────┐ │ ZooKeeper 集群(3 台服务器) │ │ 记录:谁是 Leader?分区怎么分?谁活着谁挂了? │ └────────┬──────────────────────────┬──────────────┘ │ │ ┌────▼────┐ ┌────▼────┐ ┌────▼────┐ │ Broker1 │ │ Broker2 │ │ Broker3 │ └─────────┘ └─────────┘ └─────────┘
KRaft 模式(Kafka 3.3+)
Kafka 说"我自己就能管理自己",把管理员的工作内化到 Kafka 节点内部。一部分节点同时当"管理员"(Controller)+ "分拣员"(Broker),一部分节点只当"分拣员"。
┌─────────────────────────────────────────────┐ │ Kafka 集群(不需要外部管理员) │ │ │ │ ┌──────────────────┐ ┌──────────────────┐ │ │ │ Broker1 │ │ Broker2 ││ │ │ 🧠 Controller │ │ 🧠 Controller ││ │ │ 📦 Broker │ │ 📦 Broker ││ │ └──────────────────┘ └──────────────────┘ │ │ ┌──────────────────┐ ┌──────────────────┐ │ │ │ Broker3 │ │ Broker4 ││ │ │ 🧠 Controller │ │ 📦 Broker(纯干活的)││ │ │ 📦 Broker │ │ ││ │ └──────────────────┘ └──────────────────┘ │ └─────────────────────────────────────────────┘
全方位对比
| 对比维度 | ZooKeeper 模式 | KRaft 模式 |
|---|---|---|
| 架构 | 需要额外维护 ZooKeeper 集群(至少 3 台) | 内置于 Kafka,零外部依赖 |
| 学习成本 | 两套系统:ZK 的选举机制 + Kafka 的存储机制 | 一套系统,只学 Kafka |
| 部署复杂度 | 先起 ZK 集群 → 再起 Kafka 集群,有严格的启动顺序 | 一条 docker-compose up -d 搞定 |
| 监控运维 | 两套监控、两套日志、两套报警 | 一套搞定,日志更统一 |
| Controller 切换 | 依赖 ZK 选举,心跳断开可能误判,需 30 秒+ | 基于 Raft 协议,< 10 秒完成切换 |
| 元数据同步 | Controller 从 ZK 读取,再同步给 Broker(多一步) | Controller 直接从 Raft 日志同步(更直接) |
| 分区上限 | 数万分区后 ZK 成为瓶颈 | 百万级分区无压力 |
| 启动速度 | 先恢复 ZK 状态,再恢复 Kafka 状态 | 直接恢复本地 Raft 日志,快 2-3 倍 |
| 版本趋势 | Kafka 2.x 及以前 | Kafka 3.3+ 默认,未来唯一模式 |
为什么 KRaft 更快?
打个比方:
ZooKeeper 就像你每次做事都要跑去找老大汇报:
Broker:"老大,我想当 Leader" → ZK:"我看看...行,你去吧" → Broker 把结果通知其他 Broker
KRaft 就像团队内部直接投票:
Broker1:"我提名我自己" → Broker2:"同意" → Broker3:"同意" → 好,你当 Leader
少了一个中间人,少了一次转发,自然更快。而且 Raft 协议本身就是为分布式一致性设计的,比 ZK 的 Zab 协议更简洁。
你现在用的就是这个
你 docker-compose.yaml 里配的就是 KRaft 模式:
# Controller 仲裁投票者列表——只有这 3 个参与管理决策 KAFKA_CFG_CONTROLLER_QUORUM_VOTERS: "1@kafka-1:9091,2@kafka-2:9091,3@kafka-3:9091" # kafka-1/2/3 = Controller + Broker(既管事又干活) KAFKA_CFG_PROCESS_ROLES: "broker,controller" # kafka-4 = 纯 Broker(只干活不参与管理决策) KAFKA_CFG_PROCESS_ROLES: "broker"
-
kafka-1/2/3三个组成 Raft 仲裁组,投票选 Leader Controller,容忍 1 个挂掉 -
kafka-4是纯 Broker,不参与投票,只管存储消息 -
不需要任何 ZooKeeper 容器
3.2 docker-compose.yaml 关键配置
# 集群架构: # kafka-1/2/3 = Controller(管理集群) + Broker(存储消息) # kafka-4 = 纯 Broker(只干活,不管事) # 三个监听器,各司其职: # CONTROLLER://:9091 → Controller 之间的内部通信 # BROKER://:19092 → Broker 之间的数据同步 # EXTERNAL://:9092 → 你的 Spring Boot 程序连接这个端口 # 可靠性配置: # default.replication.factor = 3 (每条消息存 3 份) # min.insync.replicas = 2 (至少 2 份写入成功才确认)
3.3 启动
# 进入目录 cd 你的项目目录 # 启动集群(首次会自动拉取镜像) docker-compose up -d # 验证 4 个节点都在运行 docker ps --filter "name=kafka" # 应该看到 kafka-1, kafka-2, kafka-3, kafka-4 四个容器
四、Spring Boot 实战:亲手收发消息
4.1 项目依赖
<dependency> <groupId>org.springframework.kafka</groupId> <artifactId>spring-kafka</artifactId> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-web</artifactId> </dependency>
4.2 连接配置
spring: kafka: bootstrap-servers: 你的IP:9092,你的IP:9093,你的IP:9094,你的IP:9095 producer: key-serializer: StringSerializer # key 怎么转成字节 value-serializer: StringSerializer # value 怎么转成字节 acks: all # 等所有副本确认 consumer: key-deserializer: StringDeserializer # 字节怎么转回 key value-deserializer: StringDeserializer # 字节怎么转回 value group-id: demo-group # 消费者组 auto-offset-reset: earliest # 从头开始读
什么是序列化? Kafka 只认识字节流(010101),不认 String 也不认对象。序列化就是把你的数据变成字节,反序列化就是把字节还原回来。就像写信:序列化 = 把字写在纸上,反序列化 = 读纸上的字。
4.3 生产者:发送消息
@Service
public class KafkaProducerService {
private final KafkaTemplate<String, String> kafkaTemplate;
// Spring Boot 自动给你准备好了 KafkaTemplate,直接注入
public KafkaProducerService(KafkaTemplate<String, String> kafkaTemplate) {
this.kafkaTemplate = kafkaTemplate;
}
// 发消息:一行代码
public void send(String key, String value) {
kafkaTemplate.send("demo-topic", key, value);
}
}
4.4 消费者:接收消息
@Service
public class KafkaConsumerService {
// @KafkaListener 注解:自动监听 topic,有新消息自动调用此方法
@KafkaListener(topics = "demo-topic", groupId = "demo-group")
public void listen(String message) {
System.out.println("收到消息: " + message);
// 这里写你的业务逻辑:存数据库、发短信、更新缓存...
}
}
就这么简单!发消息一行 kafkaTemplate.send(),收消息一个 @KafkaListener 注解。
4.5 动手测试
# 1. 启动应用 mvn spring-boot:run # 2. 浏览器里直接测试(最简单) http://localhost:8080/api/message/send?msg=HelloKafka http://localhost:8080/api/message/send?key=user001&msg=订单1 http://localhost:8080/api/message/send?key=user001&msg=订单2 # 3. 查看消费者收到了什么 http://localhost:8080/api/message/received
你会在控制台看到消费者打印的消息,也会在 API 返回中看到完整的消费记录。
五、进阶:用 Kafka 传输文件
传输文件跟传字符串的原理一样——只是序列化器不同。
| 字符串消息 | 文件传输 | |
|---|---|---|
| value 类型 | String |
byte[] |
| 序列化器 | StringSerializer |
ByteArraySerializer |
| 反序列化器 | StringDeserializer |
ByteArrayDeserializer |
本质上就是:文件 → 读成字节数组 → Kafka → 字节数组 → 写回文件。
生产者
@Service
public class KafkaFileProducerService {
private final KafkaTemplate<String, byte[]> kafkaTemplate;
public void sendFile(String fileName, byte[] fileData) {
kafkaTemplate.send("file-topic", fileName, fileData);
}
}
消费者
@Service
public class KafkaFileConsumerService {
@KafkaListener(topics = "file-topic", groupId = "file-group",
containerFactory = "byteArrayContainerFactory")
public void receiveFile(ConsumerRecord<String, byte[]> record) {
String fileName = record.key(); // 文件名
byte[] data = record.value(); // 文件内容
Files.write(Paths.get("./received-files/" + fileName), data);
}
}
测试
# 上传一张图片 curl -X POST http://localhost:8080/api/file/upload -F "file=@C:\你的图片.jpg" # 查看接收记录 curl http://localhost:8080/api/file/received # 文件保存在 ./received-files/ 目录
注意事项
-
Kafka 单条消息默认最大 1MB,大文件需要改配置或只传文件路径
-
生产环境建议传文件 OSS 地址,消费者自己去下载
六、常见问题 FAQ
Q1: Kafka 和 RabbitMQ 有什么区别?
| Kafka | RabbitMQ | |
|---|---|---|
| 设计目标 | 高吞吐的流式数据处理 | 灵活的消息路由 |
| 消息留存 | 消息持久化,可重复消费 | 消费完就删了 |
| 吞吐量 | 百万条/秒 | 万条/秒级别 |
| 适用场景 | 日志收集、实时计算、事件驱动 | 业务解耦、任务分发、RPC |
简单记:Kafka 像水库(蓄水+放水),RabbitMQ 像快递(点到点送达)。
Q2: 消息会丢吗?
Kafka 提供了三层保障:
-
生产者 acks=all:等所有副本写入成功才确认
-
分区副本:每个分区多份拷贝,Broker 挂了不影响
-
min.insync.replicas=2:至少 2 个副本确认才写入
这三层配齐,基本不会丢消息。
Q3: 消息会重复消费吗?
有可能。Kafka 保证"至少一次"投递。你的消费者要做好幂等处理(同一条消息处理两次结果一样),比如用数据库唯一键去重。
Q4: 消息顺序能保证吗?
-
同一个分区内:严格有序
-
不同分区之间:无序
所以如果你要顺序,就把需要有序的消息用相同的 key 发送,它们会进入同一分区。
Q5: 消费者组有什么用?
假设一个 Topic 有 3 个分区:
Consumer Group A: Consumer Group B: 消费者1 → Partition0 消费者4 → Partition0 消费者2 → Partition1 消费者5 → Partition1,2 消费者3 → Partition2
-
同一组内:自动分工,一条消息只被组内一个人处理(负载均衡)
-
不同组之间:互不影响,同一条消息两个组都能收到(广播)
七、学习路线图
第1步 ✅ 理解核心概念 ├── Topic / Partition / Offset ├── Producer / Consumer / Consumer Group └── Broker / Replica 第2步 ✅ 动手部署 ├── Docker Compose 启动集群 └── 理解监听器、副本、KRaft 配置 第3步 ✅ Spring Boot 收发消息 ├── 字符串消息(理解序列化) └── 文件传输(理解不同序列化器) 第4步(自行探索) ├── 发送 JSON 对象消息 ├── 手动提交 Offset(精确控制消费进度) ├── Kafka Streams(流处理) └── 配置监控(Kafka UI / Prometheus)
八、总结
回顾一下你掌握了什么:
-
为什么用 Kafka:解耦、削峰、异步,系统之间通过消息通信而不是直接调用
-
五大核心概念:Topic / Partition / Offset / Producer / Consumer Group
-
Docker 部署:4 节点 KRaft 集群,不用 ZooKeeper
-
Spring Boot 集成:生产者用
KafkaTemplate.send(),消费者用@KafkaListener -
字符串 vs 文件:区别在于序列化器(StringSerializer vs ByteArraySerializer),其余一模一样
Kafka 学起来最难的其实就是前半个小时的概念理解。一旦跑通了第一条消息,后面的都是在这条路上越走越远。你现在已经站在起点了。
更多推荐




所有评论(0)