Kafka架构与存储原理
学习 Kafka 不能只背“Topic、Partition、Consumer Group”。真正要理解的是:Kafka 为什么能高吞吐、为什么能横向扩展、为什么能让不同消费者独立回放数据。
整体架构
flowchart TD
A["Producer"] --> B["Broker 1"]
A --> C["Broker 2"]
A --> D["Broker 3"]
B --> E["Topic: order.created"]
C --> E
D --> E
E --> F["Partition 0"]
E --> G["Partition 1"]
E --> H["Partition 2"]
F --> I["Consumer Group A"]
G --> I
H --> I
F --> J["Consumer Group B"]
G --> J
H --> JKafka 集群由多个 Broker 组成。Topic 被拆成多个 Partition,Partition 分散在不同 Broker 上。生产者写 Topic 时,实际写的是某个 Partition;消费者读 Topic 时,也是从一个个 Partition 拉取数据。
Broker
Broker 是 Kafka 服务端节点。它负责:
- 接收生产者写入请求。
- 把消息追加到本地日志文件。
- 为消费者提供拉取接口。
- 复制其他 Broker 上的分区数据。
- 参与分区 Leader 选举和元数据管理。
单个 Broker 宕机会影响它负责的分区 Leader,但如果副本配置合理,其他 Broker 可以接管。
Topic 和 Partition
Topic 是业务维度的消息分类,例如:
order.created:订单创建事件。payment.success:支付成功事件。user.behavior:用户行为日志。
Partition 是 Topic 的物理分片。一个 Topic 可以有多个 Partition,每个 Partition 是一条有序日志。
flowchart TD
A["Topic: order.created"] --> B["Partition 0"]
A --> C["Partition 1"]
A --> D["Partition 2"]
B --> E["Offset 0 -> Offset 1 -> Offset 2"]
C --> F["Offset 0 -> Offset 1 -> Offset 2"]
D --> G["Offset 0 -> Offset 1 -> Offset 2"]Kafka 只能保证同一个 Partition 内消息有序,不能保证整个 Topic 全局有序。因为多个 Partition 是并行写入、并行消费的。
为什么要有 Partition
如果一个 Topic 只有一个 Partition,所有消息只能写到一个日志文件里,消费也只能被一个消费组成员处理。吞吐和扩展能力都会被单点限制。
有了 Partition 之后:
- 写入可以分散到多个 Broker。
- 消费可以由多个消费者并行处理。
- 单个 Topic 可以承载更大吞吐。
代价是:顺序只能在分区内保证。业务要根据 key 把同一实体的消息打到同一个分区,比如同一个订单号、同一个用户 ID。
Record 的结构
一条 Kafka 消息通常包括:
| 字段 | 作用 |
|---|---|
| key | 决定分区,也可用于业务幂等 |
| value | 消息内容 |
| timestamp | 事件时间或写入时间 |
| headers | 扩展元数据,例如 traceId、来源系统 |
| offset | 写入分区后的唯一位置 |
key 不是必须的,但生产环境通常建议设置。没有 key 时,消息可能被均匀分散到不同分区;有 key 时,同一个 key 通常会进入同一个分区,从而保证局部顺序。
Kafka 为什么写得快
Kafka 高吞吐主要来自几个设计:
- 顺序追加写:消息追加到日志末尾,减少随机写。
- Page Cache:依赖操作系统缓存提升读写性能。
- 批量发送:Producer 可以把多条消息合并成一个批次。
- 零拷贝:Broker 发送文件数据给消费者时可以减少用户态和内核态复制。
- 分区并行:多个分区可以在多个 Broker 上并行读写。
flowchart TD
A["Producer 批量发送"] --> B["Broker 顺序追加日志"]
B --> C["操作系统 Page Cache"]
C --> D["磁盘刷写"]
C --> E["Consumer 拉取"]
E --> F["零拷贝发送数据"]不要把“写磁盘”简单理解为“一定慢”。随机小写入确实慢,但顺序写、批量写、利用缓存后,磁盘吞吐可以非常可观。
日志段 Segment
Partition 在磁盘上不是一个无限大的文件,而是由多个 Segment 组成。每个 Segment 通常包含:
.log:真正存消息内容。.index:offset 到文件位置的索引。.timeindex:时间戳索引。
flowchart TD
A["Partition 0"] --> B["00000000000000000000.log"]
A --> C["00000000000000100000.log"]
A --> D["00000000000000200000.log"]
B --> E["index 和 timeindex"]
C --> F["index 和 timeindex"]
D --> G["index 和 timeindex"]Segment 的好处是方便删除旧数据。Kafka 根据保留策略删除整个旧 Segment,而不是逐条删除消息。这也是它高吞吐的重要原因。
副本机制
为了避免某个 Broker 宕机导致数据不可用,Kafka 给 Partition 配置副本。副本分为 Leader 和 Follower:
- Leader:处理生产者写入和消费者读取。
- Follower:从 Leader 拉取数据复制。
- ISR:保持同步的副本集合。
flowchart TD
A["Producer"] --> B["Partition 0 Leader on Broker 1"]
B --> C["Follower on Broker 2"]
B --> D["Follower on Broker 3"]
C --> E["ISR"]
D --> E
F["Broker 1 宕机"] --> G["从 ISR 中选新 Leader"]如果 Leader 宕机,Kafka 会从可用副本中选出新的 Leader。生产环境常见配置是副本因子大于 1,例如 3 副本。
LEO 和 HW
理解可靠性时,经常会遇到两个概念:
- LEO:Log End Offset,某个副本日志末尾的下一个 offset。
- HW:High Watermark,高水位,消费者能看到的最大已提交位置。
通俗讲,Leader 不能只看自己写到哪里,还要看副本有没有跟上。只有被足够副本复制到的位置,才应该对消费者可见。
flowchart TD
A["Leader 写到 Offset 10"] --> B["Follower 1 复制到 Offset 10"]
A --> C["Follower 2 复制到 Offset 8"]
B --> D["HW 推进到安全位置"]
C --> D
D --> E["Consumer 只能读取已提交消息"]如果消费者读到还没有被其他副本复制的数据,Leader 宕机后这些数据可能丢失,消费者就会读到“后来不存在”的消息。HW 的作用就是避免这种不一致。
Controller 和元数据
Kafka 需要有人管理集群元数据,例如:
- 哪些 Broker 在线。
- Topic 有哪些 Partition。
- Partition Leader 在哪个 Broker。
- Broker 宕机后如何重新选 Leader。
这个角色就是 Controller。旧部署中常见 ZooKeeper 参与元数据管理,现代部署更多使用 KRaft 模式。初学阶段不需要死背部署差异,先记住一句话:
Kafka 不只是存消息,还要维护分区、Leader、副本、消费组等元数据,否则集群不知道请求该发给谁。
如果架构理解不清会怎样
| 不理解的点 | 容易犯的错 | 后果 |
|---|---|---|
| Partition | 随便设置分区数 | 吞吐不够,或顺序被破坏 |
| key | 发送消息不带 key | 同一订单事件进入不同分区,处理顺序不可控 |
| 副本 | 副本因子设为 1 | Broker 宕机后数据不可用甚至丢失 |
| ISR | 只配置 acks=all,不配置最小同步副本 | 看似可靠,实际副本不足时仍有风险 |
| 保留策略 | 把 Kafka 当永久数据库 | 磁盘爆满,旧数据被删后无法回放 |
小结
Kafka 的核心不是“队列”,而是“分布式追加日志”。Topic 提供业务分类,Partition 提供并行和局部顺序,Replica 提供容灾,Offset 提供消费进度。把这些关系串起来,后面的生产、消费、可靠性就不再是零散配置。
