这是《Kafka 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。 本章一次性建立 Kafka 的完整概念地图。后面所有章节都会反复引用这些名词,值得花时间彻底吃透。
2.1 消息(Record)
Kafka 中数据的基本单位叫消息,也叫记录(record/event)。一条消息由以下部分组成:
┌──────────────────────────────────────┐
│ Key 可选,用于分区与业务标识 │
│ Value 消息体,业务数据本身 │
│ Timestamp 时间戳(创建时间或追加时间) │
│ Headers 可选键值对元数据 │
└──────────────────────────────────────┘
- Key 常用来标识“同一实体”,例如订单号、用户 ID。相同 key 的消息会进同一分区,从而保证该实体的消息在分区内有序;
- Value 才是真正承载业务的字节串,格式完全由业务决定(JSON、Avro、Protobuf 都行,第 8 章细讲);
- Timestamp 有两种语义:
CreateTime(生产者创建时间)和LogAppendTime(Broker 追加时间),由 topic 级参数message.timestamp.type决定; - Headers 适合放链路追踪 ID、租户标识、消息版本等元信息,不参与分区计算。
2.2 主题(Topic)与分区(Partition)
Topic 是逻辑上的消息类别,比如 order-events、app-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 并行与扩展的基本单位:
- 每个分区是一个只追加的提交日志(append-only log),消息写入后不可修改;
- 分区可以分布在不同 Broker 上,写入与读取并行度随集群规模增长;
- 分区内有序,跨分区不保证全局有序——这是 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 保存若干分区副本,负责处理读写请求;
- 每个分区有多个副本,其中一个叫 Leader,其余是 Follower;
- 所有读写都由 Leader 处理,Follower 只是异步拉取同步(第 11 章);
- 集群中有若干 Controller(KRaft 模式下是一个 Raft 组),负责元数据管理与 Leader 选举(第 12 章)。
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 Replica:处理所有读写请求;
- Follower Replica:向 Leader 拉取数据保持同步,不对外服务;
- ISR(In-Sync Replicas):与 Leader 保持同步的副本集合(含 Leader)。
当 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 │
└─────────────────────┘
一次完整的数据流:
- Producer 按 key 计算目标分区,把批次发给分区 Leader;
- Leader 追加到本地日志,Follower 拉取同步,满足
acks=all与min.insync.replicas后返回成功; - Consumer Group 的成员向协调器加入组,分区被分配给组内消费者;
- 消费者向分区 Leader 发 fetch 请求,按 offset 拉取数据;
- 处理完成后提交位移到内部主题
__consumer_offsets; - 数据按保留策略保留,与是否被消费无关。
2.8 消息保留与日志压缩
Kafka 不会在消息被消费后删除它,删除只由保留策略驱动:
retention.ms:按时间保留,默认 7 天;retention.bytes:按大小保留,默认 -1(不限);cleanup.policy=compact:日志压缩,每个 key 只保留最新一条,适合状态类数据(配置表、用户资料、CDC 最新快照);- 两者可以组合为
compact,delete。
“消费不删除”带来了两个重要能力:新消费者可以从头回放历史;同一份数据可被多个独立系统重复消费。
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 | 消费进度落后于最新位移 |
本章小结
- Topic 是逻辑分类,Partition 是物理并行单位,分区内有序、跨分区无序;
- Offset 是分区内的位置编号,消费位移与消息位移是两回事;
- 读写都走分区 Leader,Follower 只做同步; Controller 负责元数据与选举;
- 消费组内一个分区只给一个消费者,组间互不影响;
- 数据按保留策略删除,消费本身不删除,因此可回放、可多订阅。
思考题
- 要保证“同一个用户的所有行为按顺序处理”,应该怎么设计 key 和分区?
- 一个 topic 6 个分区、消费组 4 个消费者,分配结果可能是怎样的?如果消费者加到 8 个呢?
- 为什么 Kafka 的删除策略和消费进度解耦?这带来了哪些好处和代价?