读完这篇博客,你将能够:用 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 提供了三层保障:

  1. 生产者 acks=all:等所有副本写入成功才确认

  2. 分区副本:每个分区多份拷贝,Broker 挂了不影响

  3. 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 学起来最难的其实就是前半个小时的概念理解。一旦跑通了第一条消息,后面的都是在这条路上越走越远。你现在已经站在起点了。

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐