Skip to content

Kafka架构与存储原理

学习 Kafka 不能只背“Topic、Partition、Consumer Group”。真正要理解的是:Kafka 为什么能高吞吐、为什么能横向扩展、为什么能让不同消费者独立回放数据。

整体架构

mermaid
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 --> J

Kafka 集群由多个 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 是一条有序日志。

mermaid
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 上并行读写。
mermaid
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:时间戳索引。
mermaid
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:保持同步的副本集合。
mermaid
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 不能只看自己写到哪里,还要看副本有没有跟上。只有被足够副本复制到的位置,才应该对消费者可见。

mermaid
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同一订单事件进入不同分区,处理顺序不可控
副本副本因子设为 1Broker 宕机后数据不可用甚至丢失
ISR只配置 acks=all,不配置最小同步副本看似可靠,实际副本不足时仍有风险
保留策略把 Kafka 当永久数据库磁盘爆满,旧数据被删后无法回放

小结

Kafka 的核心不是“队列”,而是“分布式追加日志”。Topic 提供业务分类,Partition 提供并行和局部顺序,Replica 提供容灾,Offset 提供消费进度。把这些关系串起来,后面的生产、消费、可靠性就不再是零散配置。