Kafka入门
前言
在微服务和分布式系统日益普及的今天,消息队列对后端开发者而言基本上算是必修课了。而 Kafka,这个曾经由 LinkedIn 创造、如今属于Apache的消息系统,凭借其极致的吞吐量、持久化存储和消息回溯能力,几乎成了大数据、日志采集、实时流处理场景下的标配。
然而,Kafka 目前的学习并不是特别简单。市面上的大多数教程都比较古老,并且学习起来非常吃力,这也是我自己尝试写这篇文章的出发点。在学习中我逐渐了解到,不同于RabbitMQ的分区即是日志,同时也接触到了raft协议,更深入了解了线程的一些内容。
这份笔记,正是我学习 Kafka 的全过程记录与整理。
这份笔记,是进入AI时代之后我们普通人的一个尝试,汲取市面上各个文档教程的优点,鞭策和驱使AI,完成的这份笔记。整体笔记我是全流程学习和提问下来的,AI初次生成的版本还是非常差,而且存在各种错误,当然这篇文章目前应该还是存在不少问题,希望大家批评指正。
本文特色
-
四章入门:从“是什么”到“怎么用”,再到“Spring Boot 集成”,最后落到实战,路径算是较为清晰,一步步鞭策ai改成这个样子。
-
图示较多:关键概念基本都配有 ASCII 图示或 Mermaid 流程图,帮助我们理解内容。
-
强调“为什么”:不只列出配置项,更解释每一项的背后含义,比如为什么分区日志模型能实现高吞吐、为什么
concurrency并不会帮你自动把拉取和处理解耦到不同线程。 -
融入 KRaft 模式:当前大多教程仍基于 ZooKeeper 模式,本文专门对比了两种模式,并用 KRaft 搭建了集群环境,但是内容比较浅薄,如果需要深入原理,大家也可以去尝试鞭策AI,但一定要分辨信息。
-
可验证、可运行:每一章都设计了【验证】环节,从启动日志判断配置是否生效。
阅读建议
-
如果你是完全零基础,建议按顺序阅读即可。
-
如果你已有一定经验,我觉得这个文章对你来说帮助不大。
-
笔记中的 Demo 代码会同步到 GitHub : https://github.com/123jacj/Kafka-demo,供大家直接下载运行和扩展。
第 1 章:Kafka 概述与核心概念
本章目标:
- 理解消息队列的三大作用:解耦、削峰、异步。
- 明确 Kafka 的定位与核心能力,能和其他 MQ 对比。
- 掌握 Topic、Partition、Offset 三层结构,理解“分区即提交日志”的设计。
- 理解分区日志模型的特点(不可变、可回溯、高吞吐)及其适用场景。
- 认识 Producer、Consumer 和 Consumer Group,理解组内负载均衡与组间广播。
1.1 什么是消息队列?
消息队列(Message Queue) 是实现异步通信的关键组件。主要有两个角色:
- 生产者(Producer):发送消息的一方。
- 消费者(Consumer):接收并处理消息的一方。
消息队列的核心作用,就是连接两者实现:
- 解耦:生产者和消费者互不从属,代码功能没有任何交叉,双方互不影响。
- 削峰填谷:面对大量请求,通过队列缓冲使后端服务平稳处理,避免流量尖峰压垮系统。
- 异步处理:生产者发送消息后直接返回,不等待处理结果,提升运行效率和系统吞吐。
例如用户下单后,订单服务只需发送订单创建消息,库存、积分、短信等服务各自异步处理,无需订单服务依次调用。
1.2 什么是 Kafka?
Kafka 是一个分布式流处理平台,由 LinkedIn 开发,现为 Apache 顶级项目。它既可作为高性能消息队列,也可用于日志聚合、实时流处理和事件溯源。
核心能力
| 能力 | 说明 | 典型场景 |
|---|---|---|
| 消息队列 | 异步解耦生产者和消费者 | 订单完成后触发库存扣减 |
| 日志系统 | 收集和聚合应用日志 | 运维监控、ELK 日志分析 |
| 流处理 | 实时数据流处理 | 实时计算、实时报警 |
| 事件溯源 | 记录所有业务事件 | 审计日志、数据同步 |
Kafka vs 传统 MQ
| 特性 | RabbitMQ | ActiveMQ | RocketMQ | Kafka |
|---|---|---|---|---|
| 消息模型 | 队列 / 交换机 | 队列 | 主题 / 队列 | 分区日志 |
| 吞吐量 | ~10K/s | ~1K/s | ~100K/s | ~100K/s |
| 消息持久化 | 可选 | 可选 | 默认持久化 | 默认持久化 |
| 消息顺序 | 单队列有序 | 单队列有序 | 队列内有序 | Partition 内有序 |
| 消息保留 | 基于消费删除 | 基于消费删除 | 基于消费删除 / 时间 | 基于时间删除 |
| 时效性 | 实时 | 实时 | 实时 | 近实时 |
| 事务消息 | 不支持 | 支持 | 支持 | 支持 |
RocketMQ 在吞吐量上与 Kafka 同级别,且原生支持事务消息、顺序消息等特性,常用于电商、金融场景。
这里推荐初学者学习消息队列时 RocketMQ以及Kafka二选一。
一句话总结:Kafka 是为高吞吐量场景设计的分布式日志系统,天然支持消息持久化和多消费者组,消息被消费后不会立即删除,而是按时间策略保留。
1.3 核心概念
1.3.1 Topic(主题)
Topic 是消息的逻辑分类,类似于数据库中的表。
┌─────────────────────────────────────────────┐
│ Topic: order │
├─────────────────────────────────────────────┤
│ Partition 0 │ Partition 1 │ Partition 2 │
│ [msg 0] │ [msg 0] │ [msg 0] │
│ [msg 1] │ [msg 1] │ [msg 1] │
│ [msg 2] │ │ [msg 2] │
└─────────────────────────────────────────────┘
特点:
- 每个 Topic 可包含多个 Partition(分区)。
- 消息以追加方式写入分区末尾。
- 消息被消费后不会删除,而是根据时间策略滚动清除。
1.3.2 Partition(分区)
Partition 是 Kafka 实现并行处理和水平扩展的基础,是消息存储的最小单元。
要理解分区,我们先了解分区的特点:
-
只能追加,不能修改:
- 生产者发送到某个分区的消息总是被追加到该分区日志文件的末尾。
- 不能修改或删除分区中间的某条消息。数据一旦写入,就是不可变的(Immutable)。
-
严格有序:
- 分区内的每一条消息都会被分配一个从 0 开始递增的序号,即偏移量 (Offset)。
- 例如一个分区里有三条消息,它们的 Offset 依次是 0, 1, 2。消费者读取时,也按这个顺序,从而保证在分区内部消息的写入和消费顺序完全一致。
-
持久的事件记录:
- 分区对应 Broker 磁盘上的一个或多个日志文件,消息被持久化存储,直到超过保留策略(如 7 天或大小上限)才会被删除。
分区的作用:
- 并行处理:多个分区可被不同消费者同时消费。
- 负载均衡:消息按 Key 散列到不同分区。若 Key 为 null,则采用轮询(Round-Robin) 策略,让消息均匀分布。
- 水平扩展:增加分区数即可提升 Topic 的吞吐能力。
- 分区有序:单个分区内消息严格有序,不同分区之间无顺序保证。
分区副本(Replica):为实现高可用,每个分区可以有多个副本。其中:
- Leader:负责处理该分区的所有读写请求。
- Follower:从 Leader 同步数据,作为备份。
- ISR(In-Sync Replicas):与 Leader 保持同步的副本集合。当 Leader 故障时,会从 ISR 中选举新的 Leader。
1.3.3 Offset(位移)
每条消息在其分区内都有一个唯一且递增的序号,称为 Offset。
(Redis 中也有类似概念用于主从复制,但含义不同。)
Partition 0:
┌─────┬─────┬─────┬─────┬─────┬─────┐
│ 0 │ 1 │ 2 │ 3 │ 4 │ 5 │ ← Offset
├─────┼─────┼─────┼─────┼─────┼─────┤
│msg 1│msg 2│msg 3│msg 4│msg 5│msg 6│
└─────┴─────┴─────┴─────┴─────┴─────┘
▲
│
Consumer 从这里开始消费
committed offset = 5
next fetch = 6
关键概念:
- committed offset:消费者已确认处理完毕的位置,支持自动或手动提交。
- next fetch:下次拉取消息时从该 offset 开始。
- 消费者通过提交 offset 记录消费进度,重启后可继续消费。
1.3.4 为什么 Kafka 的消息模型是“分区日志”?
Kafka 本质上是一个分布式提交日志。每个分区在物理上对应磁盘上的一个或多个日志文件,消息只能追加,不可修改或随机删除。消费者通过 offset 像翻书一样“拉取”消息,消费进度由消费者自己提交保存。
一个Kafka分区就是一个提交日志
根据这一点,可以理解 Kafka 的许多特性:
- 高性能:顺序写磁盘非常高效,这是 Kafka 高吞吐的关键。
- 解耦:生产者把消息放入分区,消费者读取,双方完全解耦。
- 数据回溯与重放:日志不可变,消费者可以反复消费历史消息,用于故障恢复、新应用回充状态、多目标消费等。
分区日志的优缺点:
| 优点/缺点 | 具体说明 |
|---|---|
| ✅ 消息可回溯 | 重置 offset 即可重复消费历史数据,适用于故障恢复、审计等场景 |
| ✅ 高吞吐 | 顺序写入磁盘、零拷贝技术,单机可达几十万条/秒 |
| ✅ 消息持久不丢 | 消息写入后不会因消费而删除,天然持久化 |
| ✅ 多消费者独立 | 不同消费者组各自维护 offset,互不干扰 |
| ❌ 无单条消息确认 | 按 offset 提交,粒度比传统队列粗 |
| ❌ 消息不可变 | 已写入的消息无法修改,需业务自行处理幂等或补偿 |
| ❌ 延迟略高 | 批量发送和拉取导致端到端延迟在毫秒~百毫秒级,非极低延迟 |
| ❌ 分区内顺序消费的代价 | 同一分区只能被组内一个消费者消费,若需严格顺序且高并发,须足够多的分区 |
这些特点决定了 Kafka 特别适合大数据、日志、事件流等对吞吐量和数据回溯要求高的场景,而不太适合需要精确单条确认、极低延迟的金融交易指令分发等场景。
1.3.5 Producer(生产者)
生产者是向 Kafka 发送消息的客户端。
// 伪代码
ProducerRecord<String, Order> record = new ProducerRecord<>(
"order-topic", // Topic
"order-001", // Key(用于分区路由)
new Order(...) // Value(消息内容)
);
Future<RecordMetadata> future = kafkaTemplate.send(record);
发送模式:
| 模式 | 行为 | 适用场景 |
|---|---|---|
| Fire & Forget | 发送后不等待结果 | 日志、统计分析 |
| 同步发送 | 等待发送结果返回 | 需要确认发送成功 |
| 异步发送 | 回调通知发送结果 | 高吞吐量场景 |
1.3.6 Consumer(消费者)
消费者是从 Kafka 读取消息的客户端。
@KafkaListener(topics = "order-topic", groupId = "order-service")
public void consume(Order order) {
System.out.println("收到订单: " + order);
}
核心机制:
- 消费者属于某个消费者组(Consumer Group)。
- 同一个分区只能被组内的一个消费者消费。
- 不同消费者组之间互不影响,各自都可以消费到 Topic 中的全部消息。
1.3.7 Consumer Group(消费者组)
消费者组是 Kafka 灵活性的核心。它同时实现了队列模式和发布/订阅模式。
┌────────────────────────────────────────────────────┐
│ Topic: order │
│ Partition 0 │ Partition 1 │ Partition 2 │
├────────────────────────────────────────────────────┤
│
┌──────────┴──────────┐
▼ ▼
┌─────────────────┐ ┌─────────────────┐
│ Consumer Group │ │ Consumer Group │
│ "inventory" │ │ "notification" │
│ ┌─────────────┐ │ │ ┌─────────────┐ │
│ │ consumer-1 │ │ │ │ consumer-3 │ │
│ │ (P0, P1) │ │ │ │ (P0) │ │
│ └─────────────┘ │ │ └─────────────┘ │
│ ┌─────────────┐ │ │ ┌─────────────┐ │
│ │ consumer-2 │ │ │ │ consumer-4 │ │
│ │ (P2) │ │ │ │ (P1, P2) │ │
│ └─────────────┘ │ │ └─────────────┘ │
└─────────────────┘ └─────────────────┘
│ │
▼ ▼
库存服务消费 通知服务消费
(消息队列模式) (发布订阅模式)
两种模式的实现
- 队列模式(组内负载均衡):同一消费者组内的多个消费者“瓜分”Topic 的所有分区,每条消息只被组内一个消费者处理。
- 发布/订阅模式(组间广播):不同消费者组之间完全独立,各自都会消费到 Topic 的全部消息,互不干扰。
多个消费者组的意义
同一条消息往往需要被多个独立系统处理,例如“订单已支付”消息:
- 库存系统扣减库存(组
inventory) - 积分系统增加积分(组
points) - 通知系统发送短信(组
notification)
这些系统的逻辑、消费速度、代码实现完全独立,放在不同的消费者组里即可解决“一个消息被多个系统消费”的问题,同时各系统内部还可以用多消费者并发处理来加速。
1.4 Kafka 基本术语速查表
| 术语 | 中文 | 说明 |
|---|---|---|
| Broker | 代理节点 | Kafka 集群中的一个服务节点 |
| Cluster | 集群 | 多个 Broker 组成的整体 |
| Topic | 主题 | 消息的逻辑分类 |
| Partition | 分区 | Topic 的物理分片,可并行处理 |
| Replica | 副本 | 分区的备份,用于高可用 |
| Leader | 主副本 | 负责处理读写请求的副本 |
| Follower | 从副本 | 从 Leader 同步数据的副本 |
| Producer | 生产者 | 发送消息的客户端 |
| Consumer | 消费者 | 接收消息的客户端 |
| Consumer Group | 消费者组 | 消费者的工作组,分摊消费 |
| Offset | 偏移量 | 消息在分区中的唯一序号 |
| ISR | 同步副本集 | 与 Leader 保持同步的副本集合 |
1.5 本章小结
- 消息队列:异步解耦、削峰填谷、异步处理。
- Kafka 是什么:分布式流处理平台,天生为高吞吐量日志场景设计。
- 核心概念:Topic → Partition → Offset 三层结构,消息以分区日志形式持久存储,不可变、可回溯。
- 分区日志的优缺点:吞吐高、持久可靠、可多组独立消费;但牺牲了单条精确确认、极低延迟等特性。
- 生产者/消费者:异步解耦,支持多消费者组,通过分区分配实现了队列和发布/订阅两种模式。
第 2 章:Kafka 结构与工作原理入门
本章目标
- 了解 Kafka 集群架构及组件关系。
- 对比 ZooKeeper 与 KRaft 模式,理解控制器选举与 Raft 基本流程。
- 区分分区副本 Leader 与 KRaft Controller Leader 的不同。
- 掌握副本机制:Leader/Follower/ISR,以及 Broker 宕机时的选举过程。
- 熟悉消息写入流程及
acks配置对可靠性的影响。 - 理解物理存储结构:日志段、索引、活跃段与段滚动。
- 了解消费者拉取模型、位移管理,以及数据清理的基本方式。
2.1 Kafka 集群架构
2.1.1 整体架构图
┌──────────────────────────────────────────────────────────────────┐
│ Kafka 集群 │
├──────────────────────────────────────────────────────────────────┤
│ │
│ ┌───────────────┐ ┌───────────────┐ ┌───────────────┐ │
│ │ Broker 1 │ │ Broker 2 │ │ Broker 3 │ │
│ │ ┌─────────┐ │ │ ┌─────────┐ │ │ ┌─────────┐ │ │
│ │ │ Leader │ │ │ │ Leader │ │ │ │ Leader │ │ │
│ │ │ P0 副本 │ │ │ │ P1 副本 │ │ │ │ P2 副本 │ │ │
│ │ └─────────┘ │ │ └─────────┘ │ │ └─────────┘ │ │
│ │ ┌─────────┐ │ │ ┌─────────┐ │ │ ┌─────────┐ │ │
│ │ │Follower │ │ │ │Follower │ │ │ │Follower │ │ │
│ │ │ P1 副本 │ │ │ │ P2 副本 │ │ │ │ P0 副本 │ │ │
│ │ └─────────┘ │ │ └─────────┘ │ │ └─────────┘ │ │
│ └───────────────┘ └───────────────┘ └───────────────┘ │
│ │ │ │ │
│ └─────────────────────┼─────────────────────┘ │
│ │ │
│ ▼ │
│ ┌────────────────────────────────┐ │
│ │ 集群协调服务 (ZooKeeper 或 │ │
│ │ KRaft) │ │
│ │ · 元数据存储 │ │
│ │ · 控制器选举 │ │
│ └────────────────────────────────┘ │
└──────────────────────────────────────────────────────────────────┘
2.1.2 组件说明
| 组件 | 官方定义 | 说明 |
|---|---|---|
| Broker | 一个 Kafka 服务器 | 作为集群中的一个节点,负责存储消息、处理客户端请求。 |
| Topic | 消息的逻辑分类 | 生产者发布消息到 Topic,消费者从 Topic 订阅消息。 |
| Partition | 有序、不可变的消息序列 | 一个 Topic 可被分割成多个 Partition,以实现并行处理和水平扩展。 |
| Replica | 分区副本 | 同一个 Partition 在不同 Broker 上的冗余拷贝,用于实现高可用。 |
| Controller | 控制器 | 集群中一个特殊的 Broker,负责管理分区和副本的状态,执行集群管理任务。 |
控制器是什么?
Controller 是 Kafka 集群的“管理员”,本身也是一个 Broker,但同一时刻集群中只有一个 Broker 担任 Controller。它负责监听 Broker 上下线、为分区选举 Leader 副本、管理副本分配和 ISR 状态等。Controller 通过 ZooKeeper 或 KRaft 协议选举产生,当 Controller 宕机时,其他 Broker 会立即选举出新的 Controller。
2.2 控制器(Controller)
2.2.1 控制器选举:ZooKeeper 模式 vs KRaft 模式
Kafka 集群需要“谁来协调管理”的机制,这就是控制器的由来。控制器的选举方式经历了两个阶段:
| 对比项 | ZooKeeper 模式 | KRaft 模式 |
|---|---|---|
| 版本支持 | Kafka 2.x 及更早 | Kafka 3.3+(生产可用) |
| 外部依赖 | 需要独立的 ZooKeeper 集群 | 无外部依赖,完全自包含 |
| 元数据存储 | 存储在 ZooKeeper 中 | 存储在 Kafka 内部 Topic __cluster_metadata 中 |
| 选举机制 | 多个 Broker 竞争创建 ZooKeeper 临时节点,先创建者成为 Controller | 基于 Raft 共识协议 在 Controller 节点组内选举 |
| 部署复杂度 | 高,需同时维护 Kafka 和 ZooKeeper 两套集群 | 低,仅需部署 Kafka |
| 分区可扩展性 | 受 ZooKeeper 性能限制,分区数不宜过多 | 支持百万级分区 |
什么是 Raft 协议?
Raft 是一种分布式一致性算法,核心目标是让一组节点能够可靠地就“谁是 Leader”以及“共享数据的状态”达成一致。它的工作方式是:
- 集群中的节点通过选举产生一个 Leader,其余节点为 Follower。这里的 Leader 指 KRaft Controller Leader。
- 所有数据变更请求都必须先发送给 Leader,由 Leader 将变更日志复制到所有 Follower。
- 当日志被集群中多数派节点确认后,Leader 提交该变更并通知所有节点应用。
- 如果 Leader 故障,剩余节点会自动选举出新 Leader,保证系统持续运行。
在 KRaft 模式下,Kafka 用一组专用的 Controller 节点运行 Raft 协议来管理集群元数据,完全替代了 ZooKeeper,使 Kafka 成为一个真正的自包含系统。
两个“Leader”对比
| 维度 | 分区副本 Leader | KRaft Controller Leader |
|---|---|---|
| 节点是什么 | 分区的某个副本(落在 Broker 上) | 运行 Raft 的 Controller 进程(也是 Broker 或专用节点) |
| 领导者职责 | 处理该分区的生产/消费请求 | 处理集群元数据变更,管理分区副本的 Leader 选举 |
| 选举算法 | 基于 ISR 列表,由 Controller 指定 | Raft 协议(多数派投票) |
| 影响范围 | 单个分区 | 整个集群元数据 |
| 数据流 | 生产/消费消息的数据管道 | 元数据日志(Topic 配置、分区状态等) |
2.2.2 控制器的职责
Controller 职责
├── 分区 Leader 选举 # 当 Broker 故障时,为 Leader 副本失效的分区选举新 Leader
├── 元数据管理 # 处理 Topic/Partition 的创建、删除、配置变更
├── Broker 生命周期管理 # 监听 Broker 的加入和退出,更新集群成员信息
├── 副本分配 # 为新分区或扩容副本选择 Broker,决定副本在集群中的分布
└── ISR 管理 # 维护每个分区的 ISR 集合,动态移除非同步副本
2.2.3 分区 Leader 选举:以 Broker 宕机为例
关键前提:每个分区有多个副本,分散在不同 Broker 上,但同一时刻只有一个副本被选为 Leader,由它处理该分区的所有读写请求。不同分区的 Leader 分布在不同的 Broker 上,实现负载均衡。
当某个 Broker 宕机时,其上承载的所有副本都会失效。影响分为两类:
- Leader 副本在该 Broker 上的分区:该分区暂时失去服务能力,需要 Controller 立即选举新 Leader。
- 只有 Follower 副本在该 Broker 上的分区:该分区的 Leader 在其他 Broker 上,读写服务完全不受影响。
因此,单个 Broker 宕机并不等于整个集群瘫痪,只有部分 Leader 分区需要短暂的重新选举(通常毫秒到秒级)。
选举流程
1. Controller 检测到 Broker 1 宕机。
2. 找出所有 Leader 副本位于 Broker 1 的分区。
3. 对每个受影响的分区,从 ISR 中选择一个存活的 Follower 提升为新 Leader。
4. Controller 更新集群元数据,向所有 Broker 广播新的 Leader 信息。
5. 生产者和消费者收到元数据刷新,将请求转向新 Leader。
具体示例
集群有 3 个 Broker,Topic order 有 3 个分区 (P0, P1, P2),副本因子为 3。正常分布如下:
Broker 1: P0 Leader 副本, P1 Follower 副本
Broker 2: P1 Leader 副本, P2 Follower 副本
Broker 3: P2 Leader 副本, P0 Follower 副本
Broker 1 宕机后:
| 分区 | 受影响的副本 | 影响分析 | 恢复动作 |
|---|---|---|---|
| P0 | Leader 副本在 Broker 1 上,已失效 | P0 暂时无法读写 | 将 Broker 3 上的 P0 Follower 提升为新 Leader,P0 恢复 |
| P1 | Follower 副本在 Broker 1 上,已失效 | P1 的 Leader 在 Broker 2,读写不受影响 | 将 Broker 1 上的 P1 副本从 ISR 移除,P1 继续服务 |
| P2 | 无副本在 Broker 1 上 | 完全不受影响 | 无需操作 |
关键理解:失效副本从对应分区的 ISR 中剔除,并不是移除整个 Broker。整个集群中仅 P0 经历了短暂的 Leader 重选举,P1 和 P2 服务一切正常。
2.3 分区副本机制
2.3.1 副本类型与 ISR
| 术语 | 说明 |
|---|---|
| Leader Replica | 每个分区中唯一负责处理所有读写请求的副本。生产者写入、消费者读取都经过 Leader。 |
| Follower Replica | 被动从 Leader 同步数据的副本,不处理客户端请求。作用是作为备份,在 Leader 故障时接替。 |
| ISR (In-Sync Replicas) | 与 Leader 保持同步的副本集合,包含 Leader 自身及能及时追上 Leader 数据进度的 Follower。 |
ISR 的核心作用:
ISR 是 Kafka 在数据可靠性和系统可用性之间的平衡机制:
- 当生产者配置
acks=all时,消息必须被 ISR 中所有副本确认才算发送成功,这样即使 Leader 宕机,新 Leader 也从 ISR 中选出,保证数据不丢失。 - ISR 是动态的:如果一个 Follower 副本同步滞后超过
replica.lag.time.max.ms(默认 30 秒),Controller 会将其从 ISR 中移除;当它追上进度后重新加入。 - 这种设计避免了因某个慢 Follower 而阻塞整个分区的写入,兼顾了性能与可靠性。
2.3.2 副本分布示例
Topic order,3 个分区,副本因子 3。副本在集群中的分布如下:
Broker 1 Broker 2 Broker 3
┌──────────┐ ┌──────────┐ ┌──────────┐
│ P0 副本 │ │ P1 副本 │ │ P2 副本 │
│ (Leader) │ │ (Leader) │ │ (Leader) │
├──────────┤ ├──────────┤ ├──────────┤
│ P1 副本 │ │ P2 副本 │ │ P0 副本 │
│ (Follower)│ │ (Follower)│ │ (Follower)│
└──────────┘ └──────────┘ └──────────┘
分布特点:
- 每个分区的 Leader 副本分散在不同 Broker 上,实现负载均衡。
- 每个分区有三个副本分布在不同 Broker 上,任意一个 Broker 宕机不会造成数据丢失,实现高可用。
- 同一个 Broker 上承载了多个不同分区的副本,角色各不相同。
2.3.3 数据同步流程
Follower 副本通过拉取(Fetch) 方式从 Leader 副本同步数据:
Producer Leader 副本 Follower 副本
│ │ │
│ 1. 发送消息 │ │
│──────────────────────▶│ │
│ │ 2. 写入本地日志 │
│ │ │
│ │ 3. Follower 发起 Fetch │
│ │◀───────────────────────│
│ │ 4. 返回新消息 │
│ │────────────────────────▶
│ │ 5. Follower 写入确认 │
│ │◀───────────────────────│
│ 6. 根据 acks 返回 │ │
│◀──────────────────────│ │
为什么 Follower 主动拉取而不是 Leader 推送?
- Follower 可根据自身负载和网络状况控制同步速率,避免被压垮。
- 拉取模式天然支持批量处理,效率更高。
- Leader 不需要维护每个 Follower 的推送状态,设计更简洁。
2.4 消息写入流程
2.4.1 写入步骤详解
Producer 写入消息完整流程
1. 选择目标分区
- 若指定 key,则 hash(key) % 分区数 确定分区
- 若未指定 key,采用轮询或粘性分区策略
2. 查找 Leader 位置
- Producer 从任意 Broker 获取元数据,找到目标分区的 Leader 副本所在 Broker
3. 建立连接
- 与 Leader Broker 建立 TCP 连接(可复用已有连接)
4. 批量发送消息
- Producer 将多条消息打包为一个批次(batch),一次性发送
5. Leader 写入磁盘
- Leader 副本将消息追加写入本地日志文件(顺序写,性能极高)
6. 副本同步
- Follower 副本从 Leader 拉取新消息并写入各自日志
7. 返回 ACK
- Leader 根据 acks 配置,等待足够确认后向 Producer 返回确认
8. Producer 收到 ACK,认为消息发送成功
2.4.2 acks 配置与可靠性
| acks | 说明 | 可靠性 | 性能 |
|---|---|---|---|
| 0 | Producer 发送后不等待任何确认 | 极低,消息可能丢失 | 最高 |
| 1 | Leader 写入磁盘后即返回确认 | 中,Leader 宕机且未同步时可能丢失 | 高 |
| all (-1) | 等待所有 ISR 副本确认后才返回 | 最高,需配合 min.insync.replicas 使用 |
较低 |
acks=0:
Producer ──→ 发送消息,立即返回
acks=1:
Producer ──→ Leader 写入磁盘 → 返回确认
acks=all (-1):
Producer ──→ Leader 写入 → ISR 所有副本确认 → 返回确认
2.4.3 物理存储结构
1. 目录与文件概览
每个分区在 Broker 磁盘上对应一个独立的目录,消息以日志段(Log Segment) 为单位切分存储:
logs/
└── order-topic-0/ ← 分区目录
├── 00000000000000000000.log ← 消息数据文件
├── 00000000000000000000.index ← 偏移量索引
├── 00000000000000000000.timeindex ← 时间戳索引
└── 00000000000000000000.snapshot ← 事务快照(可选)
| 文件 | 说明 |
|---|---|
.log |
实际存储消息内容,按 offset 顺序追加写入 |
.index |
稀疏索引,记录 offset → 文件物理位置的映射 |
.timeindex |
记录时间戳 → offset 的映射,支持按时间范围消费 |
2. 日志段:消息的物理容器
- 逻辑上,一个分区是无限、有序的消息队列。
- 物理上,操作系统无法高效管理无限增长的单文件,因此 Kafka 将分区切分为多个大小固定的日志段。
Partition-0 (逻辑概念)
│
├── Segment 1 (00000000000000000000.log) ← offset 0 ~ 1000
├── Segment 2 (00000000000000001001.log) ← offset 1001 ~ 2000
└── Segment 3 (00000000000000002001.log) ← offset 2001 ~ ... (活跃段)
3. 段文件命名
段文件名由该段中第一条消息的 offset(Base Offset)左补零至 20 位组成。例如 00000000000000002001.log 表示首条消息 offset 为 2001。Kafka 通过二分查找文件名快速定位消息所属段。
4. 段的内部组成与索引
| 索引文件 | 作用 |
|---|---|
.index(偏移量索引) |
稀疏记录 offset → .log 文件中的物理字节位置,实现快速随机读取 |
.timeindex(时间戳索引) |
记录时间戳 → offset,支持按时间戳查找消息 |
消息查找流程:
- 定位段:根据目标 offset 二分查找段文件名。
-
索引定位:在段的
.index文件中找到最近且 ≤ 目标的索引条目,获取物理位置。 -
顺序扫描:从
.log文件该位置开始顺序读取,直到找到目标消息。
5. 活跃段与段滚动
每个分区同时只有一个段处于可写状态,称为活跃段(Active Segment)。当活跃段达到大小或时间阈值时,会被密封(变为只读),并创建新的活跃段继续写入。
6. 设计优势
- 顺序 I/O 高性能:追加写入当前活跃段,避免随机磁盘寻址。
- 高效数据管理:以段为最小物理单元,方便文件级操作。
- 快速查找:段文件命名 + 稀疏索引,O(log N) 级别定位。
- 不可变性保证:旧段只读,确保日志安全性与一致性。
2.5 消息读取流程
2.5.1 读取步骤
Consumer 读取消息流程
1. 调用 poll(),指定超时时间
2. 根据已提交的 offset 决定各分区从何处开始消费
3. 向各分配分区的 Leader 副本发送 Fetch 请求
4. Leader 利用 .index 快速定位 offset,从 .log 读取消息
5. 将消息批次返回给 Consumer
6. Consumer 处理消息后,提交 offset 记录消费进度(自动或手动)
Kafka 采用消费者拉取(Pull) 模型,消费者可根据自身处理能力控制消费速度,避免被消息推送压垮。
2.5.2 消费者位移管理
Offset 记录了消费者在每个分区上的消费进度:
分区消息序列:
[0] [1] [2] [3] [4] [5] [6] [7] [8] ...
已消费并提交 下次拉取位置
committed = 4 next fetch = 5
两种提交方式:
| 方式 | 行为 | 优缺点 |
|---|---|---|
自动提交 (enable.auto.commit=true) |
每隔 5 秒自动提交上一次 poll 的所有消息 offset | 简单,但可能重复消费(提交前宕机)或丢失消息(处理完未提交) |
手动提交 (enable.auto.commit=false) |
业务代码显式调用 commitSync() 或 commitAsync()
|
可精确控制提交时机,推荐用于要求可靠性的场景 |
手动提交的具体实现将在第 4 章 Spring Boot 集成中演示。
2.6 数据清理机制
2.6.1 保留策略
Kafka 的消息不会因为被消费而删除,而是根据预设的保留策略进行清理:
| 策略 | 说明 | 配置项 |
|---|---|---|
| 基于时间 | 超过指定时间的消息会被删除 |
retention.ms(默认 7 天) |
| 基于大小 | 分区日志总大小超过阈值时,删除最旧的消息 | retention.bytes |
两种策略可同时设置,满足任一条件即触发清理。
2.6.2 清理方式
| 方式 | 行为 | 适用场景 |
|---|---|---|
| delete | 直接删除过期消息(默认) | 常规消息流,如日志收集、用户行为埋点 |
| compact | 保留相同 Key 的最新值,删除旧版本 | 数据库变更日志、状态快照、事件溯源 |
Compact 策略示意:
Before compaction: After compaction:
Key A: value1 ─┐
Key A: value2 ├── 只保留最新版本 Key A: value3
Key A: value3 ─┘
Key B: value2 ─── 保留 Key B: value2
Key C: value1 ─── 保留(无更新) Key C: value1
这种策略非常适合将 Kafka 用作可靠的数据源,可以随时恢复每个 Key 的最新状态。
2.7 本章小结
- 集群架构:多个 Broker 组成集群,一个 Controller 负责协调管理。
- 控制器:KRaft 模式基于 Raft 协议选举,去除了 ZooKeeper 依赖。
- 分区 Leader 选举:Broker 宕机时只有 Leader 分区重新选举,服务快速恢复。
- 副本机制:Leader-Follower 模式,ISR 动态管理同步副本集合。
- 消息写入:Producer → 分区路由 → Leader 磁盘追加 → Follower 拉取同步 → ACK 返回。
- 消息读取:Consumer 主动 Pull,通过 offset 管理消费进度。
- 数据清理:基于时间或大小的保留策略,支持 delete 和 compact 两种方式。
第 3 章:环境搭建与初步感受
本章目标:
使用 Docker Compose 部署单节点和三节点 Kafka 集群,掌握 Topic 管理、消息收发、消费者组负载均衡验证等核心命令行技能。
3.1 环境准备
3.1.1 单节点部署(Bitnami 镜像,学习用)
# docker-compose.yml
version: '3.8'
services:
kafka:
image: bitnami/kafka:3.6
container_name: kafka
ports:
- "9092:9092"
environment:
- KAFKA_CFG_NODE_ID=1
- KAFKA_CFG_PROCESS_ROLES=controller,broker
- KAFKA_CFG_CONTROLLER_QUORUM_VOTERS=1@kafka:9093
- KAFKA_CFG_LISTENERS=PLAINTEXT://:9092,CONTROLLER://:9093
- KAFKA_CFG_ADVERTISED_LISTENERS=PLAINTEXT://localhost:9092
- KAFKA_CFG_LISTENER_SECURITY_PROTOCOL_MAP=CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT
- KAFKA_CFG_CONTROLLER_LISTENER_NAMES=CONTROLLER
- KAFKA_CFG_AUTO_CREATE_TOPICS_ENABLE=true
volumes:
- kafka-data:/bitnami/kafka
volumes:
kafka-data:
启动:docker-compose up -d
进入容器:docker exec -it kafka bash
Bitnami 镜像中的 Kafka 命令带 .sh 后缀(如 kafka-topics.sh)。
3.1.2 三节点集群部署(Confluent 镜像,模拟生产)
该配置已在 CentOS 7.9 + Docker 20.10 环境验证通过。
这里建议参考Kafka入门集群搭建-学习指南-CSDN博客
关键点:
- 集群 ID 必须是 Base64 编码的 UUID(使用
kafka-storage random-uuid生成)。 -
KAFKA_ADVERTISED_LISTENERS必须使用容器服务名(如kafka1:9092),禁止使用localhost,否则容器间通信会失败。
生成集群 ID:
docker run --rm confluentinc/cp-kafka:7.5.0 kafka-storage random-uuid
# 记录输出,例如:dQOUz7b6Rx6kqFJqXf9y3w
创建 docker-compose-cluster.yml(替换 CLUSTER_ID):
version: '3'
services:
kafka1:
image: confluentinc/cp-kafka:7.5.0
container_name: kafka1
ports:
- "9092:9092"
environment:
KAFKA_NODE_ID: 1
KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: 'CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT'
KAFKA_ADVERTISED_LISTENERS: 'PLAINTEXT://kafka1:9092'
KAFKA_LISTENERS: 'PLAINTEXT://0.0.0.0:9092,CONTROLLER://0.0.0.0:9093'
KAFKA_PROCESS_ROLES: 'broker,controller'
KAFKA_CONTROLLER_QUORUM_VOTERS: '1@kafka1:9093,2@kafka2:9093,3@kafka3:9093'
KAFKA_CONTROLLER_LISTENER_NAMES: 'CONTROLLER'
KAFKA_INTER_BROKER_LISTENER_NAME: 'PLAINTEXT'
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 3
KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 3
KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 2
CLUSTER_ID: '你的Base64-UUID'
volumes:
- ./data/kafka1:/var/lib/kafka/data
restart: unless-stopped
kafka2:
image: confluentinc/cp-kafka:7.5.0
container_name: kafka2
ports:
- "9094:9094"
environment:
KAFKA_NODE_ID: 2
KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: 'CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT'
KAFKA_ADVERTISED_LISTENERS: 'PLAINTEXT://kafka2:9094'
KAFKA_LISTENERS: 'PLAINTEXT://0.0.0.0:9094,CONTROLLER://0.0.0.0:9093'
KAFKA_PROCESS_ROLES: 'broker,controller'
KAFKA_CONTROLLER_QUORUM_VOTERS: '1@kafka1:9093,2@kafka2:9093,3@kafka3:9093'
KAFKA_CONTROLLER_LISTENER_NAMES: 'CONTROLLER'
KAFKA_INTER_BROKER_LISTENER_NAME: 'PLAINTEXT'
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 3
KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 3
KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 2
CLUSTER_ID: '你的Base64-UUID'
volumes:
- ./data/kafka2:/var/lib/kafka/data
restart: unless-stopped
kafka3:
image: confluentinc/cp-kafka:7.5.0
container_name: kafka3
ports:
- "9095:9095"
environment:
KAFKA_NODE_ID: 3
KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: 'CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT'
KAFKA_ADVERTISED_LISTENERS: 'PLAINTEXT://kafka3:9095'
KAFKA_LISTENERS: 'PLAINTEXT://0.0.0.0:9095,CONTROLLER://0.0.0.0:9093'
KAFKA_PROCESS_ROLES: 'broker,controller'
KAFKA_CONTROLLER_QUORUM_VOTERS: '1@kafka1:9093,2@kafka2:9093,3@kafka3:9093'
KAFKA_CONTROLLER_LISTENER_NAMES: 'CONTROLLER'
KAFKA_INTER_BROKER_LISTENER_NAME: 'PLAINTEXT'
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 3
KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 3
KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 2
CLUSTER_ID: '你的Base64-UUID'
volumes:
- ./data/kafka3:/var/lib/kafka/data
restart: unless-stopped
启动集群:
mkdir -p ./data/kafka1 ./data/kafka2 ./data/kafka3
chown -R 1000:1000 ./data # Confluent 镜像以 appuser(UID=1000) 运行
docker-compose -f docker-compose-cluster.yml up -d
验证集群:
# 查看 KRaft 仲裁状态
docker exec -it kafka1 kafka-metadata-quorum --bootstrap-server localhost:9092 describe --status
# 预期:LeaderId 不为空,CurrentVoters: [1,2,3]
# 查看元数据复制状态
docker exec -it kafka1 kafka-metadata-quorum --bootstrap-server localhost:9092 describe --replication
# 预期:所有节点 Lag 为 0
注意:Confluent 镜像的 Kafka 命令不带
.sh(如kafka-topics),Bitnami 镜像带.sh。下面示例均基于 Confluent 集群,若使用 Bitnami 单节点请自行补全.sh后缀。
3.2 Topic 操作
3.2.1 创建 Topic
docker exec -it kafka1 kafka-topics --create \
--topic order-topic \
--bootstrap-server localhost:9092 \
--partitions 6 \
--replication-factor 3
常用可选配置:
--config retention.ms=86400000 \
--config max.message.bytes=1048576
3.2.2 查看 Topic
# 列出所有 Topic
docker exec -it kafka1 kafka-topics --bootstrap-server localhost:9092 --list
# 查看指定 Topic 的详细信息
docker exec -it kafka1 kafka-topics --bootstrap-server localhost:9092 \
--describe --topic order-topic
输出解读:
Partition: 0 Leader: 1 Replicas: 1,2,3 Isr: 1,2,3
-
Leader:该分区 Leader 副本所在的 Broker ID -
Replicas:所有副本所在的 Broker 列表 -
Isr:与 Leader 保持同步的副本列表
3.2.3 修改 / 删除 Topic
# 增加分区(只能增加,不能减少)
docker exec -it kafka1 kafka-topics --bootstrap-server localhost:9092 \
--alter --topic order-topic --partitions 8
# 修改配置(例如消息保留时间改为 2 天)
docker exec -it kafka1 kafka-topics --bootstrap-server localhost:9092 \
--alter --topic order-topic --config retention.ms=172800000
# 删除 Topic
docker exec -it kafka1 kafka-topics --bootstrap-server localhost:9092 \
--delete --topic order-topic
3.3 生产者操作
3.3.1 交互式发送消息
docker exec -it kafka1 kafka-console-producer \
--bootstrap-server localhost:9092 \
--topic order-topic
3.3.2 带 Key 的消息
docker exec -it kafka1 kafka-console-producer \
--bootstrap-server localhost:9092 \
--topic order-topic \
--property parse.key=true \
--property key.separator=:
输入格式:key:value,相同 Key 的消息会进入同一分区。
3.4 消费者操作
3.4.1 消费消息
# 从头消费所有历史消息
docker exec -it kafka1 kafka-console-consumer \
--bootstrap-server localhost:9092 \
--topic order-topic \
--from-beginning
# 只消费新到达的消息
docker exec -it kafka1 kafka-console-consumer \
--bootstrap-server localhost:9092 \
--topic order-topic
3.4.2 消费者组管理
# 列出所有消费者组
docker exec -it kafka1 kafka-consumer-groups --bootstrap-server localhost:9092 --list
# 查看消费者组的消费进度
docker exec -it kafka1 kafka-consumer-groups --bootstrap-server localhost:9092 \
--describe --group inventory
输出中 LAG 表示积压量。
3.4.3 重置 Offset
# 重置到最早(执行前需停掉该组所有消费者)
docker exec -it kafka1 kafka-consumer-groups --bootstrap-server localhost:9092 \
--group inventory --reset-offsets --topic order-topic --to-earliest --execute
# 重置到最新
docker exec -it kafka1 kafka-consumer-groups --bootstrap-server localhost:9092 \
--group inventory --reset-offsets --topic order-topic --to-latest --execute
3.5 实战:消费者组队列模式验证
本实验演示同一消费者组内两个消费者负载均衡。
常见陷阱:少量消息可能因粘性分区策略全部落在部分分区,导致一个消费者无输出。
步骤:
- 确保
order-topic有 6 个分区。 - 打开多个终端。
终端 1:启动消费者 1(组 inventory)
docker exec -it kafka1 kafka-console-consumer --bootstrap-server localhost:9092 --topic order-topic --group inventory
终端 2:启动消费者 2(同一组)
docker exec -it kafka1 kafka-console-consumer --bootstrap-server localhost:9092 --topic order-topic --group inventory
终端 3:检查再均衡是否完成
docker exec -it kafka1 kafka-consumer-groups --bootstrap-server localhost:9092 --describe --group inventory
确认 CONSUMER-ID 列出现两个不同的 ID,且分区被分配给不同消费者。
终端 4:启动生产者,连续发送至少 20 条消息
docker exec -it kafka1 kafka-console-producer --bootstrap-server localhost:9092 --topic order-topic
依次输入 msg1, msg2, …, msg20。
观察:终端 1 和终端 2 都会出现消息,且一条消息只会出现在其中一个终端。
调试技巧:
- 发送消息时使用
--property parse.key=true并给不同 Key。 - 重新启动消费者时带上
--from-beginning。
效果说明:
- 同一组内消费者分摊所有分区,实现队列模式。
- 不同组可以各自消费全量消息,实现发布/订阅模式。
3.6 命令速查表
以下命令基于 Confluent 集群环境(无 .sh 后缀),Bitnami 环境请添加 .sh。
Topic 管理
| 操作 | 命令 |
|---|---|
| 创建 Topic | kafka-topics --create --topic <name> --partitions <n> --replication-factor <r> [--config key=value] |
| 列出所有 Topic | kafka-topics --list |
| 查看 Topic 详情 | kafka-topics --describe --topic <name> |
| 增加分区数 | kafka-topics --alter --topic <name> --partitions <n> |
| 修改 Topic 配置 | kafka-topics --alter --topic <name> --config <key>=<value> |
| 删除 Topic | kafka-topics --delete --topic <name> |
生产者
| 操作 | 命令 |
|---|---|
| 交互式发送(无 Key) | kafka-console-producer --bootstrap-server <host> --topic <name> |
| 带 Key 发送 | kafka-console-producer --bootstrap-server <host> --topic <name> --property parse.key=true --property key.separator=: |
| 高级用法 | 使用 --producer-property compression.type=gzip 等参数 |
消费者
| 操作 | 命令 |
|---|---|
| 从头消费 | kafka-console-consumer --bootstrap-server <host> --topic <name> --from-beginning |
| 只消费新消息 | kafka-console-consumer --bootstrap-server <host> --topic <name> |
| 指定消费者组 | kafka-console-consumer ... --group <group-name> |
| 查看消费者组列表 | kafka-consumer-groups --list |
| 查看消费者组详情 | kafka-consumer-groups --describe --group <group-name> |
| 重置 offset 到最早 | kafka-consumer-groups --reset-offsets --topic <name> --group <group-name> --to-earliest --execute |
| 重置 offset 到最新 | kafka-consumer-groups --reset-offsets --topic <name> --group <group-name> --to-latest --execute |
| 重置到指定 offset | kafka-consumer-groups --reset-offsets --topic <name>:<partition> --to-offset <offset> --execute |
集群与元数据
| 操作 | 命令 |
|---|---|
| 查看 KRaft 仲裁状态 | kafka-metadata-quorum --bootstrap-server <host> describe --status |
| 查看元数据复制状态 | kafka-metadata-quorum --bootstrap-server <host> describe --replication |
| 查看 Broker API 版本 | kafka-broker-api-versions --bootstrap-server <host> |
| 查看日志/内部 Topic |
kafka-run-class kafka.tools.DumpLogSegments --files <log-file>(调试用) |
3.7 本章小结
- 掌握了单节点和集群环境的搭建,理解了 KRaft 模式和
advertised.listeners的正确配置。 - 熟练使用
kafka-topics管理 Topic,能解读分区和 ISR 信息。 - 能够通过命令行进行消息生产与消费,理解
--from-beginning和消费者组的概念。
第 4 章:Spring Boot 集成 Kafka
4.0 本章导学
学习目标
┌────────────────────────────────────────────────────────────────────┐
│ 四阶段学习路径 │
├────────────────────────────────────────────────────────────────────┤
│ │
│ 阶段 1:跑通链路 │
│ └── 目标:快速验证 Kafka 连通性 │
│ │
│ 阶段 2:理解配置 │
│ └── 目标:理解每个配置的作用 │
│ │
│ 阶段 3:掌握核心 ⭐⭐⭐⭐⭐ │
│ └── 目标:掌握面试必问、生产必备的核心知识点 │
│ │
│ 阶段 4:扩展应用 │
│ └── 目标:实际使用,能改代码 │
│ │
└────────────────────────────────────────────────────────────────────┘
⭐ 核心知识点
┌────────────────────────────────────────────────────────────────────┐
│ 本章最核心:concurrency │
├────────────────────────────────────────────────────────────────────┤
│ │
│ concurrency 是 Spring Kafka 的灵魂: │
│ │
│ ├── 串联所有知识点 │
│ ├── 解决实际问题(poll 阻塞) │
│ ├── 面试高频问题 │
│ └── 生产环境性能关键 │
│ │
│ 学习重点: │
│ 1. 为什么需要 concurrency?(poll 单线程困境) │
│ 2. concurrency 如何工作?(多实例并行) │
│ 3. 如何验证 concurrency 生效?(启动日志) │
│ 4. concurrency 与分区的关系? │
│ │
└────────────────────────────────────────────────────────────────────┘
Demo 代码与笔记对应表
| 功能 | 代码文件 | 关键行 | 笔记章节 |
|---|---|---|---|
| ⭐ 生产者工厂 | KafkaConfig.java | 114-139 | 4.4.3 |
| ⭐ 消费者工厂 | KafkaConfig.java | 197-232 | 4.5.3 |
| ⭐ concurrency | KafkaConfig.java | 332 | 4.8.4 |
| ⭐ DLT 配置 | KafkaConfig.java | 337-343 | 4.10.3 |
| Topic 创建 | KafkaConfig.java | 66-99 | 4.6.1 |
| ⭐ JSON 生产者 | OrderProducer.java | 68-96 | 4.4.6 |
| ⭐ JSON 消费者 | OrderConsumer.java | 59-86 | 4.5.4 |
| DLT 消费 | DeadLetterConsumer.java | 全文 | 4.10.5 |
| 异步发送 | OrderProducer.java | 68-96 | 4.11.1 |
| 同步发送 | OrderProducer.java | 117-144 | 4.11.2 |
| REST 接口 | KafkaController.java | 全文 | 4.12 |
消息的一生
┌────────────────────────────────────────────────────────────────────┐
│ 订单消息的一生 │
├────────────────────────────────────────────────────────────────────┤
│ │
│ 用户请求 ──→ Controller ──→ 生产者 ──→ Kafka ──→ 消费者 ──→ 业务 │
│ 构造消息 发送消息 存储消息 拉取消息 处理 │
│ │
│ ┌──────────────────────────────────────────────────────────────┐ │
│ │ 步骤 1: 用户发起 POST /api/kafka/orders 请求 │ │
│ │ 步骤 2: Controller 构造 OrderMessage 对象 │ │
│ │ 步骤 3: OrderProducer.sendOrder() 发送到 Kafka │ │
│ │ 步骤 4: Kafka 集群存储消息(分区 + 副本) │ │
│ │ 步骤 5: OrderConsumer 拉取消息并反序列化 │ │
│ │ 步骤 6: 业务处理,手动提交 offset │ │
│ └──────────────────────────────────────────────────────────────┘ │
│ │
└────────────────────────────────────────────────────────────────────┘
【阶段 1】跑通链路
目标:快速验证 Kafka 连通性
4.1 快速开始
4.1.1 依赖
<dependencies>
<dependency>
<groupId>org.springframework.kafka</groupId>
<artifactId>spring-kafka</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
<dependency>
<groupId>org.projectlombok</groupId>
<artifactId>lombok</artifactId>
<optional>true</optional>
</dependency>
</dependencies>
4.1.2 配置
spring:
kafka:
bootstrap-servers: localhost:9092
producer:
key-serializer: org.apache.kafka.common.serialization.StringSerializer
value-serializer: org.apache.kafka.common.serialization.StringSerializer
consumer:
key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
value-deserializer: org.apache.kafka.common.serialization.StringDeserializer
group-id: test-group
auto-offset-reset: earliest
4.1.3 启动应用
./mvnw spring-boot:run
成功日志示例:
INFO --- [ main] com.example.kafka.KafkaDemoApplication : Started KafkaDemoApplication
INFO --- [ntainer#0-0-C-1] o.s.k.l.KafkaMessageListenerContainer : test-group: partitions assigned: [test-topic-0]
4.2 简单生产者与消费者
4.2.1 最简生产者
@Service
public class SimpleProducer {
@Autowired
private KafkaTemplate<String, String> kafkaTemplate;
public void send(String message) {
kafkaTemplate.send("test-topic", message);
}
}
4.2.2 最简消费者
@Service
public class SimpleConsumer {
@KafkaListener(topics = "test-topic", groupId = "test-group")
public void consume(String message) {
System.out.println("收到消息: " + message);
}
}
4.2.3 测试接口
@RestController
@RequestMapping("/api/test")
public class TestController {
@Autowired
private SimpleProducer producer;
@GetMapping("/send/{message}")
public String send(@PathVariable String message) {
producer.send(message);
return "消息已发送: " + message;
}
}
4.2.4 【验证】链路通畅了吗?
验证步骤:
1. 启动应用
2. 访问 http://localhost:8080/api/test/send/HelloKafka
3. 观察控制台日志
✅ 成功:看到 "收到消息: HelloKafka"
❌ 失败:检查 Kafka 是否启动,端口是否正确
【阶段 2】理解配置
目标:理解每个配置的作用
4.3 消息实体定义
@Data
@Builder
@NoArgsConstructor
@AllArgsConstructor
public class OrderMessage implements Serializable {
private String orderId;
private String productName;
private Integer quantity;
private BigDecimal amount;
private String status;
private String userId;
private String address;
@JsonFormat(pattern = "yyyy-MM-dd HH:mm:ss")
private LocalDateTime createTime;
@JsonFormat(pattern = "yyyy-MM-dd HH:mm:ss")
private LocalDateTime updateTime;
}
4.4 生产者配置详解
4.4.1 完整配置示例
spring:
kafka:
bootstrap-servers: 192.168.100.168:9092,192.168.100.168:9094,192.168.100.168:9095
producer:
key-serializer: org.apache.kafka.common.serialization.StringSerializer
value-serializer: org.springframework.kafka.support.serializer.JsonSerializer
acks: all
retries: 3
properties:
enable.idempotence: true
max.in.flight.requests.per.connection: 5
batch-size: 16384
linger.ms: 10
4.4.2 配置详解
| 配置项 | 说明 | 推荐值 |
|---|---|---|
bootstrap-servers |
Kafka 集群地址 | 多节点逗号分隔 |
acks |
确认级别 |
all(强可靠) |
retries |
重试次数 | 3 |
enable.idempotence |
幂等性 | true |
batch-size |
批次大小 | 16384 |
linger.ms |
批次等待 | 10ms |
4.4.3 生产者工厂配置
@Bean
public ProducerFactory<String, OrderMessage> orderProducerFactory() {
Map<String, Object> config = new HashMap<>();
config.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
config.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
config.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class);
config.put(ProducerConfig.ACKS_CONFIG, "all");
config.put(ProducerConfig.RETRIES_CONFIG, 3);
config.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true);
return new DefaultKafkaProducerFactory<>(config);
}
@Bean
public KafkaTemplate<String, OrderMessage> orderKafkaTemplate() {
return new KafkaTemplate<>(orderProducerFactory());
}
4.4.4 acks 配置说明
acks=0:发送即成功,不等待(最快,可能丢失)
acks=1:Leader 确认即可(平衡)
acks=all:所有 ISR 确认(最可靠,推荐生产环境)
4.4.5 JSON 生产者实现
@Service
public class OrderProducer {
@Autowired
private KafkaTemplate<String, OrderMessage> kafkaTemplate;
private static final String TOPIC = "order-topic";
public void sendOrderAsync(OrderMessage order) {
kafkaTemplate.send(TOPIC, order.getOrderId(), order)
.whenComplete((result, ex) -> {
if (ex != null) {
log.error("❌ 订单发送失败: {}", ex.getMessage());
} else {
log.info("✅ 订单发送成功: partition={}, offset={}",
result.getRecordMetadata().partition(),
result.getRecordMetadata().offset());
}
});
}
}
4.5 消费者配置详解
4.5.1 完整配置示例
spring:
kafka:
consumer:
key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
value-deserializer: org.springframework.kafka.support.serializer.JsonDeserializer
group-id: order-service
auto-offset-reset: earliest
enable-auto-commit: false
properties:
spring.json.trusted.packages: "*"
4.5.2 配置详解
| 配置项 | 说明 | ⚠️ 注意事项 |
|---|---|---|
group-id |
消费者组 ID | 同组内负载均衡,组间独立消费 |
auto-offset-reset |
初始消费位置 |
earliest=从头,latest=只消费新消息 |
enable-auto-commit |
自动提交 | 生产环境建议false(手动提交) |
trusted.packages |
可信包名 | ❌未配置会导致反序列化失败 |
4.5.3 消费者工厂配置
@Bean
public ConsumerFactory<String, OrderMessage> orderConsumerFactory() {
Map<String, Object> config = new HashMap<>();
config.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
config.put(ConsumerConfig.GROUP_ID_CONFIG, "order-service");
config.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
config.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, ErrorHandlingDeserializer.class);
config.put(ErrorHandlingDeserializer.VALUE_DESERIALIZER_CLASS, JsonDeserializer.class.getName());
config.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
config.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);
config.put(JsonDeserializer.TRUSTED_PACKAGES, "*");
return new DefaultKafkaConsumerFactory<>(config);
}
4.5.4 JSON 消费者实现
@Service
public class OrderConsumer {
@KafkaListener(
topics = "order-topic",
groupId = "order-service",
containerFactory = "orderKafkaListenerContainerFactory"
)
public void consumeOrder(ConsumerRecord<String, OrderMessage> record,
Acknowledgment ack) {
OrderMessage order = record.value();
log.info("📥 收到订单消息: orderId={}, product={}", order.getOrderId(), order.getProductName());
try {
processOrder(order);
ack.acknowledge();
log.info("✅ 订单处理成功");
} catch (Exception e) {
log.error("❌ 订单处理失败: {}", e.getMessage());
throw e; // 不提交 offset,让消息重新消费
}
}
}
4.5.5 ⚠️ 常见坑点:trusted.packages
问题:反序列化失败,报 ClassNotFoundException
原因:Spring 不知道可以反序列化哪些类
解决:配置 spring.json.trusted.packages: "*"
4.6 Topic 与分区
4.6.1 Topic 创建配置
@Bean
public NewTopic orderTopic() {
return TopicBuilder.name("order-topic")
.partitions(6)
.replicas(1)
.build();
}
4.6.2 配置参数说明
| 参数 | 作用 | 开发 | 生产 |
|---|---|---|---|
partitions |
并行度上限 | 3 | 9~30 |
replicas |
数据冗余 | 1 | 3 |
4.6.3 【验证】启动日志解读
INFO --- [ntainer#1-0-C-1] : order-service: partitions assigned: [order-topic-0, order-topic-1]
INFO --- [ntainer#1-1-C-1] : order-service: partitions assigned: [order-topic-2, order-topic-3]
INFO --- [ntainer#1-2-C-1] : order-service: partitions assigned: [order-topic-4, order-topic-5]
| 日志字段 | 含义 |
|---|---|
ntainer#1 |
消费者组 1(order-service) |
-0-C-1 |
Container 1, Consumer 1 |
partitions assigned |
分区分配结果 |
说明:3 个 Container(concurrency=3)各分配 2 个分区,消费者正常工作。
【阶段 3】掌握核心 ⭐⭐⭐⭐⭐
目标:掌握面试必问、生产必备的核心知识点
4.7 手动提交与可靠性
4.7.1 为什么需要手动提交?
自动提交:简单,但可能消息处理失败却已提交 offset → 丢失
手动提交:处理成功后再提交 offset → 可靠,推荐生产环境
4.7.2 AckMode 详解
| AckMode | 说明 | 使用场景 |
|---|---|---|
MANUAL |
手动提交,必须调用 ack | ⭐ 推荐生产环境 |
MANUAL_IMMEDIATE |
立即提交,不等待 | 需要精确控制 |
BATCH |
每批提交 | 批量处理 |
TIME |
定时提交 | 超时自动提交 |
COUNT |
达到数量提交 | 计数自动提交 |
4.7.3 手动提交示例
@KafkaListener(topics = "order-topic", groupId = "order-service")
public void consumeOrder(ConsumerRecord<String, OrderMessage> record,
Acknowledgment ack) {
try {
processOrder(record.value());
ack.acknowledge(); // ✅ 成功才提交
} catch (Exception e) {
throw e; // ❌ 不提交,下次重新消费
}
}
4.7.4 【验证】消息丢失了吗?
验证步骤:
1. 制造错误:修改 processOrder() 抛出异常
2. 发送消息
3. 观察日志
✅ 正确行为:消息重试,offset 未提交,重试失败后进入 DLT
❌ 错误行为:offset 已提交,消息丢失
4.8 concurrency 核心原理 ⭐⭐⭐⭐⭐
这是本章最核心的知识点!
4.8.1 问题:poll() 单线程困境
while (true) {
var records = consumer.poll(Duration.ofMillis(100));
for (var record : records) {
process(record); // 处理慢 → 阻塞 poll() → 消息积压
}
}
KafkaConsumer.poll() 是单线程事件驱动模型:
- poll() 一次调用拉取多个分区的消息
- 如果处理慢,poll() 必须等处理完才能拉下一批
- 示例:处理一条消息需 1 秒,每秒来 100 条 → 严重积压
4.8.2 解决:多实例并行
┌─────────────────────────────────────────────────────────────┐
│ concurrency 解决方案:一个实例一根线程 │
├─────────────────────────────────────────────────────────────┤
│ Container-1 (线程-1) │
│ ┌───────────────────────────────────────┐ │
│ │ KafkaConsumer-1 │ │
│ │ ├─ poll() 拉取 partition-0,1 │ │
│ │ ├─ 调用 @KafkaListener 处理 │ │
│ │ └─ ack / 提交 offset │ │
│ └───────────────────────────────────────┘ │
│ │
│ Container-2 (线程-2) │
│ ┌───────────────────────────────────────┐ │
│ │ KafkaConsumer-2 │ │
│ │ ├─ poll() 拉取 partition-2,3 │ │
│ │ ├─ 调用 @KafkaListener 处理 │ │
│ │ └─ ack / 提交 offset │ │
│ └───────────────────────────────────────┘ │
│ │
│ Container-3 (线程-3) │
│ ┌───────────────────────────────────────┐ │
│ │ KafkaConsumer-3 │ │
│ │ ├─ poll() 拉取 partition-4,5 │ │
│ │ ├─ 调用 @KafkaListener 处理 │ │
│ │ └─ ack / 提交 offset │ │
│ └───────────────────────────────────────┘ │
│ │
│ 关键: │
│ - concurrency=3 会创建 3 个 KafkaMessageListenerContainer │
│ - 每个 Container 在自己的线程里跑 │
│ - 该线程负责 poll、反序列化、调用监听器、提交 offset │
│ - 拉取与处理都在**同一个线程**内串行执行 │
│ - 通过增加实例数,让不同分区的消息能被并行处理 │
└─────────────────────────────────────────────────────────────┘
每个消费者实例的全部操作都在自己的线程内完成,且实例和线程一旦创建就会一直存在,直到应用关闭。
4.8.3 线程工厂运行机制 ⭐
ConcurrentKafkaListenerContainerFactory
│
└── concurrency=3 → 创建 3 个 KafkaMessageListenerContainer
│
├── Container-1 → 线程-1
│ └── KafkaConsumer-1
│ ├── poll() 拉取 partition-0,1
│ ├── 反序列化
│ ├── 调用 @KafkaListener
│ └── ack / 提交 offset
│
├── Container-2 → 线程-2
│ └── KafkaConsumer-2 (类似)
│
└── Container-3 → 线程-3
└── KafkaConsumer-3 (类似)
关键点:
- 每个 Container 拥有自己的线程,线程内串行执行 poll → 处理 → 提交
- 没有默认的“处理线程池”,拉取和处理共用同一线程
- 如果监听器内业务很慢,会直接阻塞当前 consumer 的 poll
- concurrency > 分区数时,多余的 consumer 会空闲
工作流程:
-
poll()拉取一批消息 - 反序列化
- 调用
@KafkaListener方法 - 根据 AckMode 提交 offset
- 回到开始等待,继续下一轮 poll
线程-1:poll → 处理 → 提交 → poll → 处理 → 提交 → ...
线程-2:poll → 处理 → 提交 → poll → 处理 → 提交 → ...
线程-3:poll → 处理 → 提交 → poll → 处理 → 提交 → ...
因此,concurrency 解决高吞吐的方式是“让不同分区的消息能被多个线程并行处理”,而非“拉取和处理解耦到不同线程池”。
4.8.4 Demo 中的 concurrency 配置
@Bean
public ConcurrentKafkaListenerContainerFactory<String, OrderMessage>
orderKafkaListenerContainerFactory() {
ConcurrentKafkaListenerContainerFactory<String, OrderMessage> factory =
new ConcurrentKafkaListenerContainerFactory<>();
factory.setConsumerFactory(orderConsumerFactory());
factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL);
factory.setConcurrency(3); // ⭐ 创建 3 个消费者实例,每个实例独立线程
DefaultErrorHandler errorHandler = new DefaultErrorHandler(
new DeadLetterPublishingRecoverer(orderKafkaTemplate()),
new FixedBackOff(1000L, 3L)
);
factory.setCommonErrorHandler(errorHandler);
return factory;
}
4.8.5 如何验证 concurrency 生效?⭐⭐⭐
启动日志:
INFO --- [ntainer#1-0-C-1] : order-service: partitions assigned: [order-topic-0, order-topic-1]
INFO --- [ntainer#1-1-C-1] : order-service: partitions assigned: [order-topic-2, order-topic-3]
INFO --- [ntainer#1-2-C-1] : order-service: partitions assigned: [order-topic-4, order-topic-5]
| 日志字段 | 含义 |
|---|---|
ntainer#1 |
消费者组 1(order-service) |
-0-C-1 |
Container 1, Consumer 1 |
partitions assigned |
分区分配结果 |
验证:启动应用后应看到 3 行日志(与 concurrency 一致),每个 Container 分配 2 个分区(6 分区 / 3)。
4.8.6 concurrency 与分区的关系
配置:partitions=6, concurrency=3
Container-1 → KafkaConsumer-1 → 消费 partition-0,1
Container-2 → KafkaConsumer-2 → 消费 partition-2,3
Container-3 → KafkaConsumer-3 → 消费 partition-4,5
关键点:
- 1 个 KafkaConsumer 可以消费多个分区
- concurrency 控制消费者数量,不是分区数
- concurrency > 分区数时,多余的消费者空闲
- 正确配置:concurrency <= 分区数
4.8.7 【验证】concurrency 生效了吗?
验证步骤:
1. 启动应用
2. 观察日志中的 "partitions assigned"
3. 确认有 N 行(N = concurrency 配置值)
✅ 成功:3 行日志,每行分配 2 个分区
❌ 失败:只有 1 行或 0 行
4.8.8 配置建议
| 场景 | 推荐 concurrency | 原因 |
|---|---|---|
| 业务处理快(打印日志) | 1 | 单线程足够 |
| IO 密集(DB、RPC) | 分区数 | 每个分区只能由一个 consumer 消费,超出分区数的线程会空闲,与分区数持平即可 |
| CPU 密集(复杂计算) | CPU 核心数 | 减少线程切换开销,且仍应 ≤ 分区数 |
| 分区数多(>10) | 分区数 | 一个线程处理一个分区 |
4.8.9 多个 @KafkaListener 的线程隔离 ⭐
每个 @KafkaListener 在启动时都会创建一组独立的、常驻的消费者线程,不同 Listener 的线程组彼此隔离。
Spring 应用启动
│
├── @KafkaListener(topic = "order-topic", groupId = "order-service")
│ ├── 线程-1 (Consumer-1) → poll order-topic → 调用 orderHandler()
│ ├── 线程-2 (Consumer-2) → poll order-topic → 调用 orderHandler()
│ └── 线程-3 (Consumer-3) → poll order-topic → 调用 orderHandler()
│
├── @KafkaListener(topic = "payment-topic", groupId = "payment-service")
│ ├── 线程-4 (Consumer-4) → poll payment-topic → 调用 paymentHandler()
│ └── 线程-5 (Consumer-5) → poll payment-topic → 调用 paymentHandler()
│
└── @KafkaListener(topic = "sms-topic", groupId = "sms-service")
└── 线程-6 (Consumer-6) → poll sms-topic → 调用 smsHandler()
隔离边界:
| 隔离维度 | 说明 |
|---|---|
| Topic | 每个 @KafkaListener 订阅自己的 Topic,不会串消息 |
| 消费者组 | 不同业务的 groupId 独立,消费进度各自维护 |
| 线程 | 每个 @KafkaListener 拥有一组专属常驻线程,互不抢占 |
| 配置 | 不同业务可使用不同的 containerFactory,独立设置并发、提交模式等 |
典型配置:
// 订单服务:concurrency=3,手动提交,启用 DLT
@KafkaListener(topics = "order-topic", groupId = "order-service", containerFactory = "orderKafkaListenerContainerFactory")
// 短信服务:concurrency=2,手动提交
@KafkaListener(topics = "sms-topic", groupId = "sms-service", containerFactory = "smsKafkaListenerContainerFactory")
// 日志服务:concurrency=1,自动提交
@KafkaListener(topics = "log-topic", groupId = "log-service", containerFactory = "logKafkaListenerContainerFactory")
每个
@KafkaListener从生到死只服务于自己绑定的 “Topic + 消费者组 + 处理逻辑”,互不干扰。
4.9 消费者组机制
4.9.1 消费者组概念
Producer ──→ order-topic
│
├──→ order-service(扣减库存)
└──→ notification-service(发送通知)
- 同组内负载均衡,不同组间独立消费
4.9.2 【验证】分区分配正确吗?
启动应用后观察 order-service 的分区分配:
- Container-1: partition-0,1
- Container-2: partition-2,3
- Container-3: partition-4,5
4.10 错误处理与死信队列(DLT)
4.10.1 重试机制
DefaultErrorHandler errorHandler = new DefaultErrorHandler(
new DeadLetterPublishingRecoverer(orderKafkaTemplate()),
new FixedBackOff(1000L, 3L) // 间隔 1 秒,重试 3 次
);
4.10.2 DLT 原理
消费消息 → 业务处理失败 → 重试 → 超过重试次数 → 发送到 DLT
DLT 消费者进行人工干预或补偿
4.10.3 Demo 中的 DLT 配置
DefaultErrorHandler errorHandler = new DefaultErrorHandler(
new DeadLetterPublishingRecoverer(orderKafkaTemplate()),
new FixedBackOff(1000L, 3L)
);
factory.setCommonErrorHandler(errorHandler);
4.10.4 DLT Topic 创建
@Bean
public NewTopic deadLetterTopic() {
return TopicBuilder.name("order-topic.DLT")
.partitions(3)
.replicas(1)
.build();
}
4.10.5 DLT 消费与排查
@KafkaListener(topics = "order-topic.DLT", groupId = "dlq-processor")
public void consumeDeadLetter(ConsumerRecord<String, OrderMessage> record,
Acknowledgment ack) {
log.warn("⚠️ 收到死信消息: orderId={}", record.value().getOrderId());
handleDeadLetter(record);
ack.acknowledge();
}
4.10.6 DLT 消息排查
通过消息头可获取原始 topic、分区、offset、异常信息。
4.10.7 【验证】DLT 生效了吗?
制造连续失败 → 重试 3 次 → 进入 order-topic.DLT → DLT 消费者收到消息
【阶段 4】扩展应用
4.11 异步与同步发送
4.11.1 异步发送(推荐)
public void sendOrderAsync(OrderMessage order) {
kafkaTemplate.send(ORDER_TOPIC, order.getOrderId(), order)
.whenComplete((result, ex) -> { ... });
}
4.11.2 同步发送
public SendResult<String, OrderMessage> sendOrderSync(OrderMessage order) {
return kafkaTemplate.send(ORDER_TOPIC, order.getOrderId(), order).get();
}
4.11.3 带分区的发送
kafkaTemplate.send(ORDER_TOPIC, partition, order.getOrderId(), order);
4.11.4 发送 API 汇总
kafkaTemplate.send("topic", "message"); // 无 Key
kafkaTemplate.send("topic", "key", "message"); // 有 Key
kafkaTemplate.send("topic", 0, "key", "message"); // 指定分区
kafkaTemplate.send("topic", "key", "message").whenComplete(...); // 异步回调
kafkaTemplate.send("topic", "key", "message").get(); // 同步阻塞
4.12 REST 接口封装
@RestController
@RequestMapping("/api/kafka")
public class KafkaController {
@PostMapping("/orders")
public ResponseEntity<ApiResponse<OrderMessage>> createOrder(
@RequestParam String productName,
@RequestParam Integer quantity,
@RequestParam BigDecimal amount,
@RequestParam String userId) {
OrderMessage order = OrderMessage.builder()
.orderId(UUID.randomUUID().toString())
.productName(productName)
.quantity(quantity)
.amount(amount)
.status("CREATED")
.userId(userId)
.createTime(LocalDateTime.now())
.build();
orderProducer.sendOrderAsync(order);
return ResponseEntity.ok(ApiResponse.success("订单已提交", order));
}
@PostMapping("/test")
public ResponseEntity<ApiResponse<String>> sendTestMessage(
@RequestParam(required = false) String key,
@RequestParam(defaultValue = "Hello Kafka!") String value) {
testProducer.sendTestMessage(key, value);
return ResponseEntity.ok(ApiResponse.success("消息已发送", value));
}
}
测试命令:
curl -X POST "http://localhost:8080/api/kafka/orders?productName=iPhone15&quantity=1&amount=7999.00"
4.13 API 文档集成(可选)
依赖:
<dependency>
<groupId>org.springdoc</groupId>
<artifactId>springdoc-openapi-starter-webmvc-ui</artifactId>
<version>2.3.0</version>
</dependency>
访问地址:
| 地址 | 说明 | 推荐 |
|---|---|---|
| http://localhost:8080/doc.html | Knife4j UI(中文) | ⭐⭐⭐⭐⭐ |
| http://localhost:8080/swagger-ui.html | Swagger UI | ⭐⭐⭐⭐ |
| http://localhost:8080/v3/api-docs | OpenAPI JSON | ⭐⭐⭐ |
【阶段 5】本章小结
4.14 本章小结
核心知识回顾:
- 阶段 1:KafkaTemplate + @KafkaListener 基本使用
- 阶段 2:bootstrap-servers、acks、group-id、trusted.packages
- 阶段 3:手动提交、concurrency(多实例并行)、消费者组、DLT
- 阶段 4:异步/同步发送、REST 接口、API 文档
生产环境 Checklist:
- replicas >= 2(推荐 3)
- min.insync.replicas = 2
- acks = all
- enable.idempotence = true
- 消费者手动提交 offset
- 配置死信队列
- concurrency <= 分区数
- 监控消费者 lag
- 业务实现幂等性
- 相同 Key 保证顺序
相关图示
1. 异常处理流程
重试循环
是
否
是
否
消费消息
业务逻辑
成功?
ack.acknowledge
提交 offset
throw e
重试 < 3?
等待 1 秒
发送到 DLT
ack 原始 offset
✅ 完成
⚠️ 死信处理
2. 发送链路
生产者内部
是
可重试异常
不可重试异常
POST /api/kafka/orders
Controller
构造 OrderMessage
KafkaTemplate.send()
ProducerFactory 获取 Producer
序列化: JsonSerializer
计算分区
KafkaProducer 发送
发送成功?
whenComplete 回调
自动重试
直接失败
3. 消费链路
错误处理
消息循环
启动与并发
是
否
是
否
Spring Boot 启动
扫描 @KafkaListener
ContainerFactory
concurrency=3
创建 3 个 Container
SimpleAsyncTaskExecutor
启动线程
poll 拉取
ConsumerFactory
反序列化: JsonDeserializer
@KafkaListener 方法
业务成功?
ack.acknowledge
提交 offset
抛出异常
重试 < 3?
FixedBackOff 等待
发送 DLT
人工处理
✅ 完成
⚠️ 死信
4. 消息流转时序图
消费侧
发送侧
ErrorHandler
@KafkaListener
Container
Kafka Broker
KafkaProducer
KafkaTemplate
Controller
用户
ErrorHandler
@KafkaListener
Container
Kafka Broker
KafkaProducer
KafkaTemplate
Controller
用户
1. 消息生产与发送
2. 消费者持续拉取
loop
[持续 poll]
3. 业务处理与容错
alt
[成功]
[重试中]
[重试耗尽]
POST 请求
1
sendOrderAsync()
2
序列化 & 分区
3
发送消息
4
ack
5
回调
6
拉取
7
消息
8
调用监听器
9
ack
10
提交 offset
11
异常
12
等待后重试
13
重新处理
14
发送 DLT
15
DLT ack
16
提交原始 offset
17
5. 完整链路
结果
消费
集群
生产
请求
成功
失败
用户
Controller
KafkaTemplate
ProducerFactory
序列化
Broker 分区
持久化
poll 拉取
反序列化
@KafkaListener
处理?
✅ ack
重试/DLT
6. 工厂与组件关系
配置注入
启动时
发送时
KafkaTemplate
ProducerFactory
KafkaProducer
序列化器
Spring 启动
ContainerFactory
Container ×3
ConsumerFactory
KafkaConsumer
Broker
application.yml
KafkaConfig.java
7. 关键配置点
容器
ContainerFactory
concurrency=3
AckMode.MANUAL
DefaultErrorHandler
FixedBackOff 重试
DLT 死信
消费者
ConsumerFactory
group-id
auto-offset-reset=earliest
JsonDeserializer
trusted.packages=*
生产者
ProducerFactory
acks=all
retries=3
enable.idempotence=true
JsonSerializer
写在最后
这份笔记确实是我个人学习Kafka相关内容的参考,并且在学习的过程中一直发现问题,提出疑问,最终形成了这版文章,其中不免存在错误和遗漏,能读到这里的小伙伴也是很厉害,希望大家能有自己的收获。
如果你在学习过程中遇到任何问题,或者发现了更好的表述方式,欢迎在评论区提出。让我们一起学习。
祝你学习愉快,编码顺利!
附录:Kafka 学习资源
| 资源 | 链接 |
|---|---|
| 官方文档 | https://kafka.apache.org/documentation/ |
| Spring Kafka | https://spring.io/projects/spring-kafka |
| 《Kafka 权威指南》 | O’Reilly 出版 |
| Apache Kafka 教程 | Apache Kafka 教程_w3cschool |
参考了 Kafka 学习指南,基于《Kafka 权威指南2》编写。
如果本文对你有帮助,欢迎点赞、收藏、关注,后续会继续更新我学到的后端和计算机知识。
再次欢迎大家批评指正,也欢迎在评论区交流你遇到的问题!