Kafka学习笔记(一):从订单事件理解核心概念与架构

订单创建成功之后,短信服务要发通知,统计服务要更新报表,推荐服务还想记录一次购买行为。

如果订单服务挨个调用这些服务,下游越多,接口越复杂。统计服务临时不可用时,难道用户也要等着?如果明天又加一个风控服务,是不是还要改订单接口?

我们可以换一种思路:订单服务只记录“订单已经创建”这件事,其他服务根据自己的需要读取这个事件。

Kafka 就是用来承接这类事件流的基础设施。

本系列以 Kafka 4.0 系列、KRaft 模式为学习基线,实操固定使用 4.0.0,方便统一环境;这个版本号是示例基线,不代表最新版本。Java 示例使用 JDK 17。

系列导航:

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

一、Kafka 是什么?

一句话理解:

Kafka 是一个能够持续写入、保存和读取事件的分布式日志平台。

这里的“日志”不只是程序打印的错误日志,还包括订单创建、支付完成、库存变化等业务事件。

Kafka 的三个基本动作是:生产者写入事件,Broker 保存事件,消费者读取事件。官方介绍将它定位为事件流平台。

1.1 为什么用“日志”来理解?

想象一本只允许往后写的账本:

1
2
3
4
位置 0:订单 O1001 已创建
位置 1:订单 O1002 已创建
位置 2:订单 O1001 已支付
位置 3:订单 O1003 已创建

短信服务读到位置 2,统计服务读到位置 1。它们各自保存一个“书签”,互不影响。

这就是后面要学习的 Offset。账本记录了什么,与某个读者读到哪里,是两件独立的事。

1.2 Kafka 能解决哪些问题?

问题 引入 Kafka 后的变化 还需要考虑什么
服务耦合 下游订阅事件,生产者无需逐个调用 事件字段与兼容性
非核心任务阻塞接口 通知、统计等任务异步处理 失败补偿与监控
瞬时流量过高 先保存事件,下游逐步处理 积压是否会超过保留期
多系统需要同一份数据 不同消费组独立读取 各组自己的消费进度
统计逻辑需要重算 在记录仍保留时重新读取 重放带来的业务副作用

消息队列让系统能够缓冲压力,但不能让持续超负荷的下游自动变快。 如果长期生产速度大于消费速度,积压仍会持续增长。


二、核心概念:先把这几个名词串起来

2.1 概念速览

概念 一句话解释 订单场景中的例子
Producer 负责写入事件的客户端 订单服务
Consumer 负责读取事件的客户端 通知服务
Broker 保存日志、处理客户端请求的服务节点 Kafka 服务实例
Topic 一类事件的逻辑名称 order-events
Partition Topic 内的一条有序日志 order-events-0
Offset 记录在某个分区中的位置 Partition 0 的位置 42
Consumer Group 一组协作读取订阅分区的消费者 notification-service
Replica 一个分区日志的副本 分布在不同 Broker 上的三份日志

把它们放到一张图里:

1
2
3
4
5
6
7
8
9
10
11
                       Topic: order-events
┌──────────────────────┐
订单服务 ── 写入 ──→ │ P0: [0][1][2][3]... │
│ P1: [0][1][2]... │
│ P2: [0][1]... │
└──────────────────────┘
│ │
▼ ▼
通知消费组 统计消费组
各自分配分区 各自分配分区
各自保存进度 各自保存进度

这里只画了逻辑结构。实际部署时,分区及其副本会分布在 Broker 上。

2.2 一条事件包含什么?

学习时可以先把一条记录理解成:

1
2
3
4
5
Topic:     order-events
Key: O1001
Value: {"eventId":"E1001","orderId":"O1001","type":"CREATED"}
Timestamp: 事件对应的时间戳
Headers: 可选的追踪信息、事件类型等

orderId 与 eventId 分工不同:

  • orderId:这是谁的事件,可以用作分区 Key。
  • eventId:这是哪一次事件,可以用于消费去重。

同一笔订单可能先创建、再支付、再退款,这三次事件应该有不同的 eventId。


三、Topic 与 Partition:分类和并行

3.1 Topic 是逻辑分类

比如按事件领域划分:

1
2
3
order-events   → 订单事件
payment-events → 支付事件
user-events → 用户事件

Topic 更接近“这一类事件放在哪里”。真正承载有序日志的是 Partition。官方概念说明介绍了这种分区组织方式。

3.2 为什么需要分区?

如果所有订单都写进同一条日志,处理并发就受到这条日志的限制。拆成多个分区后,可以将工作分散到不同节点、不同消费者。

但要记住一个边界:

Kafka 保证分区内的日志顺序,Topic 的多个分区之间没有统一的全局顺序。

例如:

1
2
P0:A1 → A2 → A3
P1:B1 → B2 → B3

读取 P0 时,日志顺序是 A1、A2、A3。至于 A1 与 B1 谁先被业务处理,Kafka 没有为它们规定一个共同顺序。

3.3 Key 如何影响顺序?

默认分区逻辑下,有 Key 的记录会根据 Key 选择分区。在分区数量和分区策略保持稳定时,相同 Key 会进入相同分区。生产者配置文档也列出了自定义分区器等选项。

订单事件可以使用 orderId 作为 Key:

1
2
3
O1001:CREATED → PAID → SHIPPED
↓
同一个 Partition

这里仍有两个前提:生产端先按正确业务顺序发送,消费端再按分区顺序执行。如果把消费结果交给线程池并发落库,执行完成的顺序仍可能被打乱。

增加分区数也可能改变 Key 到分区的映射。 对严格有序的业务,扩分区需要配合迁移或版本切换,不能只把分区数改大。

3.4 分区数与副本数有什么区别?

配置 主要解决什么问题 增加后意味着什么
分区数 数据与消费工作的并行 增加可分配的日志条数
副本数 节点故障时的数据冗余 同一条日志保存更多份

例如,3 个分区、每个分区 3 个副本,共有 9 份分区副本,但每个消费组仍只是在分配 3 个逻辑分区。


四、Offset:记录的位置与消费者的书签

4.1 Offset 只在分区内有意义

下面两条记录都可以拥有 Offset 42:

1
2
order-events / Partition 0 / Offset 42
order-events / Partition 1 / Offset 42

所以定位一条 Kafka 记录,需要组合使用 Topic、Partition、Offset。Offset 也不是业务事件的全局 ID。

4.2 消费者提交的是什么?

如果一个消费者已经处理完 Offset 0、1、2,它应该提交的进度是:

1
下一次从 Offset 3 继续

即:提交的是下一个待处理位置。 消费者当前读取位置与持久化的已提交位置还可能不同,这部分会在第三篇展开。KafkaConsumer API对两者有明确区分。

4.3 消费后消息会被删除吗?

不会因为某个消费组读取或提交 Offset 就删除。

可以把两件事分开理解:

1
2
记录保留多久 → Topic 的保留与清理策略
某组读到哪里 → 该组对各分区提交的 Offset

对于 cleanup.policy=delete,旧日志按保留时间、大小等规则清理;compact 则围绕 Key 保留状态。清理通常以日志段为单位进行,并非每条消息到时立即删除。日志实现说明解释了分段存储。

未消费不等于永久保留。 如果消费组落后太久,所需记录已经被清理,就无法凭 Offset 找回这些内容。


五、Consumer Group:组内分工,组间独立

假设 order-events 有三个分区:

1
2
3
4
5
6
notification-service
Consumer A ← P0、P1
Consumer B ← P2

analytics-service
Consumer C ← P0、P1、P2

通知服务可以用两个实例分担工作,统计服务则独立读取整份事件流。

本系列讨论的是通过 KafkaConsumer.subscribe() 使用的常规消费组。在这种模型下,稳定分配时,一个分区在同一个组里只交给一个消费者;一个消费者可以负责多个分区。KafkaConsumer 文档介绍了这种组内分配规则。

因此,一个组只订阅一个含 3 个分区的 Topic 时,即使启动 5 个消费者,也最多有 3 个消费者分到分区。

不过,“一个分区由一个消费者负责”并不意味着“业务永远只执行一次”。崩溃、重启、Rebalance 后,记录仍可能被再次处理。


六、副本与 KRaft:两层不同的职责

6.1 Leader、Follower 与 ISR

一个分区的副本可以分布在不同 Broker 上:

1
2
3
4
Partition 0
Broker 1:Leader
Broker 2:Follower
Broker 3:Follower

生产者向分区 Leader 写入,Follower 从 Leader 复制日志;普通消费通常也从 Leader 读取。ISR 是保持同步的副本集合,包含 Leader。官方副本设计说明了复制与故障接管的机制。

如果 Follower 长时间跟不上,它可能被移出 ISR。acks=all、min.insync.replicas 与 ISR 的关系,是第四篇可靠性分析的重点。

6.2 KRaft 管理的是集群元数据

Kafka 4.0 已移除 ZooKeeper 模式,使用 KRaft 管理元数据。官方升级说明明确了这一版本边界。

可以先记住这组职责:

角色 主要职责
Broker 保存业务分区日志、处理读写请求
Controller 管理集群元数据、协调分区 Leader 等状态

本地实验可以让一个进程同时担任两种角色。Controller 的元数据共识,与业务分区的 Leader/Follower 复制,需要分别理解。


七、Kafka 为什么适合高吞吐场景?

先从几个设计思路理解,不用一开始就背性能数字:

  • 顺序追加:日志持续向后写,减少随机写入。
  • 批量传输:多条记录组成批次,分摊网络和请求开销。
  • 页缓存:利用操作系统缓存服务读写。
  • 分区并行:把数据和处理工作分散出去。
  • 压缩:按批次压缩,减少传输与存储量。

这些机制在官方设计文档中有进一步解释。它们也存在权衡:批次等待更久可能增加延迟,分区更多会增加管理成本。

所以“Kafka 很快”不能替代具体测量。消息大小、Key 分布、压缩方式、副本同步、磁盘和消费逻辑都会影响最终结果。


八、学习自检

试着不看上文回答四个问题:

  1. 通知和统计都要读到订单事件,应该使用同一个组还是不同组?
  2. 一个 Topic 有 4 个分区、3 个副本,同组启动 8 个消费者,能有多少个分到分区?
  3. P0 的 Offset 10 与 P1 的 Offset 10,是同一条记录吗?
  4. 同一订单的事件使用相同 Key,消费端为什么仍要考虑执行顺序?

参考答案: 不同组;最多 4 个;不是;分区日志有序不等于异步业务执行有序。