KafkaNotes

第 02 章:核心概念全景

zjc 于 2026-01-02 发布

这是《Kafka 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。 本章一次性建立 Kafka 的完整概念地图。后面所有章节都会反复引用这些名词,值得花时间彻底吃透。

2.1 消息(Record)

Kafka 中数据的基本单位叫消息,也叫记录(record/event)。一条消息由以下部分组成:

┌──────────────────────────────────────┐
│ Key        可选,用于分区与业务标识    │
│ Value      消息体,业务数据本身        │
│ Timestamp  时间戳(创建时间或追加时间) │
│ Headers    可选键值对元数据            │
└──────────────────────────────────────┘

2.2 主题(Topic)与分区(Partition)

Topic 是逻辑上的消息类别,比如 order-eventsapp-logs。生产者向 topic 发消息,消费者订阅 topic。

每个 topic 被切分成若干个 partition

topic: order-events (3 个分区)

partition-0: [m0] [m1] [m2] [m3] ... ->
partition-1: [m0] [m1] [m2] ... ->
partition-2: [m0] [m1] [m2] [m3] [m4] ... ->

分区是 Kafka 并行与扩展的基本单位:

什么消息会进哪个分区?由生产者的分区策略决定(第 5 章):指定分区 > 按 key 哈希 > 粘性轮询。

2.3 位移(Offset)与消费位置

每个分区内的每条消息都有一个单调递增的编号,叫 offset。它由 Broker 在写入时分配,唯一标识“这条消息在该分区中的位置”。

partition-0: [0] [1] [2] [3] [4] [5] ...
                          ^
              消费者当前位置:offset = 3
              下一条要消费的是 offset = 3

关于 offset 有三组容易混淆的概念:

概念 含义
消息位移 消息在分区日志中的位置,Broker 分配,不可变
消费位移(committed offset) 消费组在某分区上“已经提交”的位置,存在 __consumer_offsets
当前位置(position) 消费者内存中下一次 poll 要读取的位置

“提交位移 100”的含义是:已经处理完 offset < 100 的所有消息,下次从 100 开始消费

2.4 Broker 与集群

一台 Kafka 服务器进程就是一个 Broker。多个 Broker 组成集群:

Broker 本身很“笨”:它不理解业务,只管按 topic/partition 存取字节流。这让它可以做到极高的通用性与吞吐。

2.5 生产者与消费者

Producer 向 topic 发布消息;Consumer 从 topic 读取消息。

Kafka 是拉模式:消费者主动调用 poll() 拉取数据,按自己的节奏消费。好处是天然背压——消费者处理不过来时,只需要降低拉取频率,不会被打垮。

消费组(Consumer Group)

消费者必须属于某个消费组(由 group.id 标识)。消费组的规则:

topic 有 3 个分区,组 A 有 2 个消费者:

  消费者1 <- partition-0, partition-1
  消费者2 <- partition-2

组 B 有 5 个消费者:

  消费者1 <- partition-0
  消费者2 <- partition-1
  消费者3 <- partition-2
  消费者4 <- 空闲
  消费者5 <- 空闲

注意最后一行:消费者数量超过分区数时,多出来的消费者空闲。想让消费更快,要么加分区,要么在下游再做并行(第 6 章、第 17 章)。

2.6 副本(Replica)与 ISR

为了容错,每个分区可以有多个副本,分布在不同 Broker 上:

当 Leader 所在 Broker 宕机,Controller 会从 ISR 中选一个 Follower 成为新 Leader,保证已提交的数据不丢失。细节(LEO、HW、leader epoch)见第 11 章。

2.7 一张全景架构图

                         ┌──────────────┐
                         │   Producers  │
                         └──────┬───────┘
                                │ produce
                                v
   ┌──────────────── Kafka Cluster ────────────────┐
   │  Broker1        Broker2        Broker3        │
   │ ┌─────────┐   ┌─────────┐   ┌─────────┐      │
   │ │ P0(L)   │   │ P0(F)   │   │ P0(F)   │      │
   │ │ P1(F)   │   │ P1(L)   │   │ P1(F)   │      │
   │ │ P2(F)   │   │ P2(F)   │   │ P2(L)   │      │
   │ └─────────┘   └─────────┘   └─────────┘      │
   │      Controller / KRaft Quorum               │
   └───────────────┬───────────────────────────────┘
                   │ fetch
                   v
        ┌─────────────────────┐
        │   Consumer Groups   │
        │  group-a   group-b  │
        └─────────────────────┘

一次完整的数据流:

  1. Producer 按 key 计算目标分区,把批次发给分区 Leader;
  2. Leader 追加到本地日志,Follower 拉取同步,满足 acks=allmin.insync.replicas 后返回成功;
  3. Consumer Group 的成员向协调器加入组,分区被分配给组内消费者;
  4. 消费者向分区 Leader 发 fetch 请求,按 offset 拉取数据;
  5. 处理完成后提交位移到内部主题 __consumer_offsets
  6. 数据按保留策略保留,与是否被消费无关。

2.8 消息保留与日志压缩

Kafka 不会在消息被消费后删除它,删除只由保留策略驱动:

“消费不删除”带来了两个重要能力:新消费者可以从头回放历史;同一份数据可被多个独立系统重复消费。

2.9 内部主题

Kafka 用普通 topic 机制实现自己的关键功能,值得认识三个内部主题:

主题 作用
__consumer_offsets 50 个分区(默认),存储消费组提交的位移
__transaction_state 存储事务状态,事务协调器使用
__cluster_metadata KRaft 模式下的元数据日志(Raft 日志)

看到这些主题不要惊讶,更不要删除它们。它们本身也是观察 Kafka 内部行为的窗口(第 12 章)。

2.10 核心术语速查表

术语 英文 一句话解释
消息 record / message / event 数据基本单位
主题 topic 逻辑消息类别
分区 partition 并行与顺序的单位,只追加日志
位移 offset 消息在分区内的编号
副本 replica 分区的拷贝,分布在不同 Broker
Leader leader replica 处理读写请求的副本
ISR in-sync replicas 与 Leader 同步的副本集合
Broker broker 一台 Kafka 服务器进程
Controller controller 管理元数据与 Leader 选举的“大脑”
生产者 producer 消息发布方
消费者 consumer 消息读取方
消费组 consumer group 组内竞争分区,组间独立广播
再平衡 rebalance 消费组内分区重新分配
位移提交 offset commit 记录消费进度的动作
消息堆积 consumer lag 消费进度落后于最新位移

本章小结

思考题

  1. 要保证“同一个用户的所有行为按顺序处理”,应该怎么设计 key 和分区?
  2. 一个 topic 6 个分区、消费组 4 个消费者,分配结果可能是怎样的?如果消费者加到 8 个呢?
  3. 为什么 Kafka 的删除策略和消费进度解耦?这带来了哪些好处和代价?