主题
第 2 章 核心概念与架构总览
学习目标:把第 1 章「Kafka 是分布式提交日志服务」这句话彻底落地 —— 能用一张大图说清楚「一条消息从生产者敲下回车到被消费者打印出来,中间经过了哪些角色、哪些通信、哪些磁盘动作」;能脱口而出 Producer / Broker / Consumer / Topic / Partition / Replica / Offset / Consumer Group / Coordinator / Controller / KRaft 这 11 个核心术语的「是什么、像什么、在哪用」;能用
kafka-topics.sh --describe在自己集群里把这些术语「肉眼看见」;最后能跑一段 Python 探针脚本,把整个集群拓扑打印成一张表格。
0. 导读:把一封信寄出去要经过多少个环节?
第 1 章我们把 Kafka 的「为什么」讲清楚了 —— 它生于 LinkedIn、底座是不可变追加日志、五大性能基石让它吞吐爆表。从本章开始,我们正式进入「是什么 / 怎么组织起来的」。
为了不让读者一上来被一堆英文术语劝退,本章会用「邮局 + 流水线 + 调度中心」这一条贯穿始终的生活化比喻:
╔════════════════════════════════════════════════╗
║ 把 Kafka 集群想象成一个「分布式邮局集团」 ║
╠════════════════════════════════════════════════╣
║ ║
║ 寄信人 Producer → 把信投到对应「柜子」 ║
║ 信件柜 Topic → 按业务分类的逻辑容器 ║
║ 流水线 Partition → 柜子内部的多条平行轨道 ║
║ 流水号 Offset → 每条流水线上消息的编号 ║
║ 分拣中心 Broker → 一台 Kafka 服务器 ║
║ 备份分拣中心 Replica → 同条流水线的多份拷贝 ║
║ 收信人 Consumer → 主动来取信的人 ║
║ 家庭信箱 Consumer Group → 同一户人家共享一个邮箱 ║
║ 班长 Coordinator → 负责给家庭成员分活的人 ║
║ 集团调度中心 Controller → 全集团人事 / 资源调度 ║
║ 集团元数据账本 KRaft → 自带 Raft 共识的元数据库 ║
║ ║
╚════════════════════════════════════════════════╝读完本章你应当能:
- 复述 11 个核心术语,且每个都讲得出「英文 → 中文 → 类比 → 在集群里的物理 / 逻辑位置」。
- 看懂一张大图:从 Producer
produce()到 Consumerpoll()中间发生了什么。 - 说出 Kafka 的「三横三纵」:存储层 / 协议层 / 客户端层 × 生产 / 消费 / 管理。
- 认识 3 个隐藏内部 Topic:
__consumer_offsets/__transaction_state/__cluster_metadata,并知道它们各自承载什么职责。 - 动手跑一遍
kafka-topics.sh --describe+ Pythoncluster_probe.py,让概念在自己屏幕上「活」起来。
1. 角色全家福(11 个核心概念,逐个击破)
每个概念都按照「首次出现 → 一句话定义 → 生活类比 → 在集群里的位置 → 关联术语」五步走。
1.1 Producer(生产者) —— 寄信人
- 首次出现:第 1 章 §1.6.3,我们调用
producer.produce(topic, key, value)把消息扔进 Kafka。 - 一句话定义:主动把消息发送到指定 Topic 的客户端进程。
- 生活类比:邮局门口的「寄信人」 —— 写好信、按格式塞进信封、丢到邮局对应的信箱柜(Topic)就走人,不关心信什么时候送到、被谁拆开。
- 物理位置:纯客户端进程,与 Broker 没有依存关系,可以是 Web 后端、IoT 设备、Connect Source、Streams 应用 …… 任何能说 Kafka 协议(TCP + 二进制)的东西都能当 Producer。
- 核心动作:
produce(topic, key, value, headers, partition?, timestamp?)—— 内部其实是「序列化 → 选分区 → 攒 batch → 网络发送 → 等 ACK」一长串。第 4 章细拆。 - 关联术语:Partitioner(决定 Key 落到哪个分区)、Acks(等几份副本确认)、Idempotence(幂等保证不重复)、Transactional Producer(跨分区事务)。
1.2 Broker(代理 / 节点) —— 分拣中心
- 首次出现:每次我们写
bootstrap.servers=localhost:9092时,连的就是一个 Broker。 - 一句话定义:一台运行 Kafka 服务进程的机器(或 Docker 容器、K8s Pod)。
- 生活类比:一个城市的「邮件分拣中心」 —— 收来自寄信人的信件,按柜子(Topic)和流水线(Partition)摆放,提供给收信人来取。集群通常 3 台起步(容忍 1 台挂掉)。
- 物理位置:有状态服务端进程,每个 Broker 维护:
- 自己负责的 Partition 副本(数据 =
log.dirs/<topic>-<partition>/*.log); - 与其他 Broker 之间的副本同步连接;
- 与 Controller 之间的元数据订阅关系;
- 对外的 9092(PLAINTEXT)/ 9093(SSL)端口。
- 自己负责的 Partition 副本(数据 =
- 关键属性:每个 Broker 有唯一
broker.id(整数),加 1 个可选的broker.rack(用于跨机架副本分布)。 - 关联术语:Listener(对外端口)、Log Dir(存储目录)、Controller(角色之一,下面会讲)、Replica(角色之一)。
📌 小白常见误区:
- 「Kafka 集群挂了」≠「所有 Broker 全挂」 —— 只要 ISR 还有机器在,分区就能继续读写。
- 「Broker」≠「Topic」 —— 一个 Broker 上同时承载几十个 Topic 的若干 Partition 副本,不是「一台机器只放一个 Topic」。
1.3 Consumer(消费者) —— 收信人
- 首次出现:第 1 章 §1.6.3,
consumer.poll(timeout)把消息拉回来。 - 一句话定义:主动从 Broker 拉取(pull)消息的客户端进程。
- 生活类比:邮局窗口前的「收信人」 —— 自己拿着家庭邮箱的钥匙(
group.id)来取信,邮局只告诉你「你家邮箱里现在有这些信」,不会主动追到你家。 - 物理位置:纯客户端进程。可以是后端服务、Streams 应用、Connect Sink、ksqlDB 查询、ETL Job ……
- 核心动作:
subscribe([topics])→poll()→ 业务处理 →commit()(提交 Offset)。第 5 章细拆。 - 关联术语:Consumer Group(多消费者协作)、Offset(消费进度)、
auto.offset.reset(首次消费的位置策略)、Rebalance(组内分区重分配)。
1.4 Topic(主题) —— 信件分类柜
- 首次出现:每次
kafka-topics.sh --create --topic learn.01.hello。 - 一句话定义:消息的逻辑分类容器,名字必须唯一。
- 生活类比:邮局大厅里贴着标签的「信件分类柜」 ——
订单事件柜、用户行为柜、支付通知柜,寄信和收信都按柜子来。 - 物理位置:逻辑概念,本身没有任何数据 —— 真正的数据在它的多个 Partition 上。Topic 在元数据里的体现就是
(topic_name, topic_id, num_partitions, replication_factor, configs[])这一行。 - 命名约定:本教程统一
learn.<chapter>.<scene>,例learn.02.cluster。生产环境推荐<env>.<domain>.<event>风格,如prod.order.created。 - 关联术语:Partition(物理切片)、Replica(副本)、Topic Config(保留策略、压缩策略、min.insync.replicas …)。
1.5 Partition(分区) —— 同类信件的多条流水线
- 首次出现:
--partitions 3创建 Topic 时指定。 - 一句话定义:Topic 内部的物理切分单元,Kafka 一切并行 / 顺序 / 副本的最小粒度。
- 生活类比:「订单事件柜」内部不是一条流水线,而是「分拣线 0 / 分拣线 1 / 分拣线 2」三条平行轨道。同一类信件按某个规则(如收件人邮编)分到不同轨道,单条轨道内严格有序,但跨轨道之间不保证。
- 物理位置:磁盘上一个独立目录
<log.dirs>/<topic>-<partition>/,里面装一组 Segment(.log/.index/.timeindex)。 - 关键性质(第 1 章 §1.3.3 已强调):
- 并行度的边界:N 个 Partition 最多被 N 个 Consumer 并行消费;
- 顺序的边界:单 Partition 内严格有序,跨 Partition 不保证;
- 副本的边界:副本以 Partition 为粒度复制;
- Rebalance 的最小单元;
- 存储的物理单元。
- 关联术语:Replica(同一 Partition 的多份拷贝)、Leader / Follower(副本角色)、ISR(同步中的副本集合)、Segment(Partition 内的日志分段)。
1.6 Replica(副本) —— 同条流水线的多份拷贝
- 首次出现:
--replication-factor 3创建 Topic 时指定。 - 一句话定义:同一个 Partition 在不同 Broker 上的多份完全相同的拷贝,用来抗 Broker 宕机。
- 生活类比:每条分拣线在 3 个不同城市都有一条「同步复刻」的备份分拣线,**正本(Leader)**接收新信、把每条信复制给「副本(Follower)」们;正本所在城市挂了,立刻让某个副本「晋升」成新正本继续干活。
- 物理位置:每个副本就是某台 Broker 的
<log.dirs>/<topic>-<partition>/目录。3 副本意味着同一份数据在 3 台机器上各占一份磁盘空间。 - 角色细分:
- Leader:唯一对外读写的副本,所有 Producer / Consumer 流量都打到它。
- Follower:只做一件事 —— 不停从 Leader 拉数据
Fetch,把自己日志追到一致;不直接对外服务。 - ISR(In-Sync Replicas):当前与 Leader 保持「足够新」的副本集合(含 Leader 自己)。掉队太久会被踢出 ISR,追上后再加回来。
- 关联术语:HW(High Watermark,已被 ISR 全部确认的最高 Offset)、LEO(Log End Offset,单个副本本地最高 Offset)、Leader Epoch(任期编号,第 9 章详讲)。
📌 副本是「整 Partition 复刻」,不是「Topic 复刻」:3 副本 × 12 分区 × 1 个 Topic = 集群里物理上有 36 份「Partition 副本」,分布在 N 个 Broker 上。
1.7 Offset(偏移量) —— 流水线上的流水号
- 首次出现:消费输出里的
partition=0 offset=5。 - 一句话定义:单个 Partition 内消息的整数编号,从 0 开始单调递增,永不重复、永不回填。
- 生活类比:分拣线 0 上每封信都贴一个递增流水号
0, 1, 2, 3 …;流水线 1 也有自己的0, 1, 2 …,两条流水线的流水号互相独立。 - 物理位置:写入 Partition 时由 Broker(Leader)分配,写到 RecordBatch 头里。Offset = (segment baseOffset + 该 batch 在 segment 中的位置 + Record 在 batch 中的偏移)。
- 关键认识:
- 生产侧
offset由 Broker 分配,Producer 不能指定。 - 消费侧 Consumer 自己维护「我读到哪儿了」,提交回
__consumer_offsets。 - Compaction Topic 的 Offset 会有「空洞」(被压缩掉的旧值),Consumer 看到的 Offset 不连续。
- 生产侧
- 三个易混 Offset:
LogStartOffset(LSO):Partition 当前最早可读的 Offset(被删除策略清理过)。HighWatermark(HW):被 ISR 全部确认、Consumer 可读的最高 Offset。LogEndOffset(LEO):单副本本地写到的最高 Offset(含未被确认的部分)。
- 关联术语:
__consumer_offsets、auto.offset.reset、seek、Consumer Lag(=LEO - CommittedOffset)。
1.8 Consumer Group(消费者组) —— 一家人共享的邮箱
- 首次出现:
group.id=learn-01配置。 - 一句话定义:一个逻辑标签(
group.id字符串),所有挂着同一个标签的 Consumer 实例共同消费一个或多个 Topic,同组内每条消息只被一人取走。 - 生活类比:「张家三口」共用一个家庭邮箱(
group.id=zhang-family),邮箱里来了一封信,三个人里任何一个取走就行;隔壁「李家」(group.id=li-family)有自己独立的邮箱,同一封广告信会发两份 —— 一家一份。 - 物理位置:逻辑概念 + 服务端协调状态。Group 元数据由集群中某台 Broker 上运行的 Group Coordinator 维护(在内存里),Offset 提交记录持久化到
__consumer_offsets这个内部 Topic。 - 核心规则:
- 同一个 Group 内:每个 Partition 只能被 1 个 Consumer 实例消费 —— 单播 / 负载均衡。
- 不同 Group 之间:每个 Partition 被每个 Group 独立消费一次 —— 广播 / 多订阅。
- 关联术语:Coordinator(协调者)、Rebalance(再平衡)、Partition Assignment Strategy(Range / RoundRobin / Sticky / Cooperative)、Static Membership(
group.instance.id)。
1.9 Coordinator(协调者) —— 给家庭成员分活的班长
⚠️ Kafka 里有两个 Coordinator,别混淆:
1.9.1 Group Coordinator(消费组协调者)
- 一句话定义:负责管理一个 Consumer Group 成员关系(谁来了 / 谁走了 / 给谁分配哪些 Partition)和 Offset 提交的 Broker。
- 生活类比:「家庭班长」 —— 张家三口都是消费成员,班长决定老大负责取奇数日邮件、老二取偶数日、老三取广告,谁请假就重新分配(Rebalance)。
- 物理位置:就是某台普通 Broker,根据公式
Math.abs(group.id.hashCode()) % numPartitions(__consumer_offsets)算出该 Group 对应的__consumer_offsets分区,该分区的 Leader 所在的 Broker 就是该 Group 的 Coordinator。 - 职责:处理
JoinGroup/SyncGroup/Heartbeat/LeaveGroup/OffsetCommit/OffsetFetch等 6 类 Request。 - 关联术语:第 11 章「Consumer Group 与 Rebalance」整章拆解。
1.9.2 Transaction Coordinator(事务协调者)
- 一句话定义:负责管理事务 Producer 状态机的 Broker,写入
__transaction_state内部 Topic。 - 生活类比:「公证处」 —— 负责签字盖章「这一组跨多个分区的写入要么全成功、要么全失败」。
- 物理位置:根据
transactional.id哈希到__transaction_state的某个分区,该分区 Leader 即 Transaction Coordinator。 - 关联术语:第 13 章「幂等与事务:Exactly Once 的真相」详讲。
1.10 Controller(集团调度中心)
- 一句话定义:整个 Kafka 集群里唯一负责管理「Broker 上下线」「分区 Leader 选举」「ISR 变更」「Topic 创建 / 删除」等集群级元数据的角色。
- 生活类比:「集团总部调度中心」 —— 哪个城市的分拣中心着火了?立刻安排别的城市顶上;新开了一个柜子(Topic),把它的分拣线分配到几个城市;某条流水线的正本掉链子,立刻晋升副本接班。
- 物理位置:
- ZK 时代:Controller 是「任意一台 Broker 兼任」,通过抢 ZK 临时节点选举出来,全集群同时只有一台。
- KRaft 时代:Controller 是「专门的 Controller 节点(Quorum)」,3-5 台组成 Raft 集群,由 Raft 选出 Active Controller,其它是 Standby。
- 关键命令:
- 看 Active Controller:
kafka-metadata-quorum.sh --bootstrap-server localhost:9092 describe --status - 老命令(仍可用):
kafka-metadata-shell.sh --snapshot ...
- 看 Active Controller:
- 关联术语:KRaft、
__cluster_metadata、Controller Epoch、Metadata Topic。
📌 超易混点:「Controller」管的是「集群拓扑 + Topic / Partition 元数据」;「Group Coordinator」管的是「消费组成员 + Offset」;「Transaction Coordinator」管的是「事务状态机」。三者都跑在 Broker 上,但职责完全不同。
1.11 KRaft(Kafka Raft 元数据模式)
- 一句话定义:Kafka 自研的、基于 Raft 共识协议的元数据管理方案,从 2.8 起预览、3.3 起生产可用、4.0 起完全替代 ZooKeeper。
- 生活类比:以前集团总部的人事档案放在「外部档案馆(ZooKeeper)」,跨部门改东西全靠档案馆挨条审;KRaft 改成「集团总部内部直接维护一本带流水的元数据账本」,一组 Controller 节点用 Raft 复制这本账本,速度更快、可用性更高、维护更简单。
- 物理位置:元数据本身存在一个特殊的内部 Topic
__cluster_metadata(注意带下划线 cluster,不是 consumer)—— 由 Controller Quorum 维护,每条元数据变更都是一条「事件」追加进去。 - 本章只做概念性介绍,第 10 章详讲:Controller Quorum、Raft 选举、Metadata Snapshot、从 ZK 迁移到 KRaft 的路径、Kafka 4.0 完全去 ZK 的影响。
- 关联术语:Controller、
__cluster_metadata、Process Roles(broker/controller/broker,controller)、Metadata Log。
📌 本教程默认环境:3 broker + 同时承载 Controller 角色的 KRaft 单机演示集群(参考根目录
docker-compose.yml)。生产推荐「3 个独立 Controller + N 个独立 Broker」分离部署。
2. 一张大图:Kafka 集群读写链路全景
2.1 Mermaid 版本(结构清晰)
2.2 ASCII 全景(信息密度更高)
┌─────────────────────────────────────────────┐
│ 客 户 端 层 │
│ ┌────────┐ ┌────────┐ ┌────────┐ ┌────────┐ │
│ │ProdA │ │ProdB │ │GroupX │ │GroupY │ │
│ │订单服务 │ │埋点上报│ │数仓ETL │ │实时风控 │ │
│ └───┬────┘ └───┬────┘ └───▲────┘ └───▲────┘ │
└─────┼──────────┼─────────┬┼──────────┼──────┘
│ Produce │ ││ Fetch │
▼ ▼ │▼ │
╔═══════════════════════════════════════════════════════════════════╗
║ KAFKA 集 群 (KRaft 模 式) ║
║ ║
║ ┌──────────────┐ ┌──────────────┐ ┌──────────────┐ ║
║ │ Broker 1 │ │ Broker 2 │ │ Broker 3 │ ║
║ │ (Standby Ctl)│ │ (Standby Ctl)│ │ Active Ctl ⭐ │ ║
║ │ ┌──────────┐ │ │ ┌──────────┐ │ │ ┌──────────┐ │ ║
║ │ │P0 Leader │◄┼──┼─┤P0 Foll. │◄┼──┼─┤P0 Foll. │ │ Topic A ║
║ │ │P1 Foll. │ │ │ │P1 Leader │ │ │ │P1 Foll. │ │ ║
║ │ │P2 Foll. │ │ │ │P2 Foll. │ │ │ │P2 Leader │ │ ║
║ │ └──────────┘ │ │ └──────────┘ │ │ └──────────┘ │ ║
║ │ Replication │◄─►│ Replication │◄─►│ Replication │ ║
║ │ Fetcher 线程 │ │ Fetcher 线程 │ │ Fetcher 线程 │ ║
║ │ │ │ │ │ │ ║
║ │ Group Coord. │ │ Tx Coord. │ │ Group Coord. │ ║
║ │ for "GroupX" │ │ for "txn-1" │ │ for "GroupY" │ ║
║ └──────┬───────┘ └──────┬───────┘ └──────┬───────┘ ║
║ │ Raft 复制 │ Raft 复制 │ Raft 复制 ║
║ └─────────┬─────────┴────────┬───────┘ ║
║ ▼ ▼ ║
║ ┌──────────────────┐ ┌──────────────────┐ ║
║ │__cluster_metadata│ │ __consumer_offsets│ ║
║ │ (KRaft Log) │ │ __transaction_state│ ║
║ └──────────────────┘ └──────────────────┘ ║
╚═══════════════════════════════════════════════════════════════════╝
▲
│ DescribeCluster / CreateTopic / DescribeConfigs
┌────┴────┐
│AdminClient│
│ (运维) │
└─────────┘2.3 一条消息的完整生命周期(时序图)
12 步完整解释:
| 步 | 动作 | 关键参数 / 配置 |
|---|---|---|
| 1 | Producer 攒一个 batch 后通过 TCP 发 ProduceRequest 给 Leader | linger.ms / batch.size / compression.type |
| 2 | Leader 把 batch 顺序追加到本地 .log,给每条 Record 分配递增 Offset | log.dirs / segment.bytes |
| 3-6 | Follower 后台 Fetcher 线程不停拉 Leader 数据,同步到本地 | replica.fetch.max.bytes / num.replica.fetchers |
| 7 | 每次 Follower Fetch 都汇报自己的 LEO;Leader 把 ISR 中所有副本的 LEO 取 min,推进 HW | replica.lag.time.max.ms |
| 8 | acks=all 时等 HW ≥ 本次写入 offset 才回 ACK;acks=1 只等 Leader;acks=0 不等 | acks / min.insync.replicas |
| 9-10 | Consumer pull 模型,每次 poll() 实际是 FetchRequest;只能读到 HW 之前的数据 | fetch.min.bytes / fetch.max.wait.ms |
| 11 | 业务代码处理消息(反序列化、入库、调下游 API 等) | —— |
| 12-13 | Consumer 把「我已经处理到 offset N+1」写到 __consumer_offsets | enable.auto.commit / auto.commit.interval.ms |
📌 关键洞察:消息的「写入路径」比「读取路径」复杂得多 —— 因为写入要保证多副本的强一致;读取只是从 Leader 单点流式读出。这就是 Kafka「写多读少不耗 CPU、读多写少看磁盘」的根因。
3. 三横三纵:Kafka 的分层架构
把整套系统横切成 3 层、纵切成 3 类职责,得到一张 3×3 的能力矩阵。这张图能帮你把第 1 ~ 21 章的所有知识点定位到具体格子里。
3.1 三横(Stack 自底向上)
┌─────────────────────────────────────────────────────────────┐
│ 三. 客户端层(Client SDK / 工具 / 生态) │
│ librdkafka / kafka-clients / Streams / Connect / ksqlDB │
│ 命令行 kafka-*.sh / Kafka UI / Cruise Control │
├─────────────────────────────────────────────────────────────┤
│ 二. 协议层(Wire Protocol & In-Broker Subsystems) │
│ Produce / Fetch / Metadata / OffsetCommit / Heartbeat ... │
│ Replication / Group Coordinator / Transaction Coordinator │
├─────────────────────────────────────────────────────────────┤
│ 一. 存储层(Disk + PageCache + 网络) │
│ Topic-Partition-Segment 三层目录 │
│ .log / .index / .timeindex / leader-epoch-checkpoint │
│ PageCache / sendfile / mmap │
└─────────────────────────────────────────────────────────────┘- 存储层(第 7、8、14 章):把日志安全、高效地写到磁盘并能快速读出来。Kafka 在这一层的工程哲学是「让 OS 替我们做缓存、让磁盘做顺序 IO、让 CPU 几乎不动」。
- 协议层(第 4、5、9、10、11、13 章):在日志之上定义「Producer / Consumer / Replica / Coordinator / Controller」之间的对话规则。每条 Kafka 协议都有明确的 ApiKey(如
Produce=0、Fetch=1、Metadata=3、OffsetCommit=8…)。 - 客户端层(第 3、12、16、17、18 章):让人类 / 业务系统能方便地说 Kafka 协议;包括所有官方 / 第三方语言客户端、命令行工具、UI、Connect Source/Sink、Streams 应用。
3.2 三纵(按职责切分)
┌────────┬────────┬────────┐
│ 生 产 │ 消 费 │ 管 理 │
┌───────────────┼────────┼────────┼────────┤
│ 客户端层 │ ProduceAPI │ ConsumeAPI │ AdminClient │
│ (用户接口) │ kafka-console-producer │ console-consumer │ kafka-topics.sh│
├───────────────┼────────┼────────┼────────┤
│ 协议层 │ ProduceReq │ FetchReq + OffsetCommit │ Metadata + DescribeConfigs│
│ (Wire Proto) │ + Idempotent + Txn │ + JoinGroup + Heartbeat │ + Reassignment等 │
├───────────────┼────────┼────────┼────────┤
│ 存储层 │ append-only .log │ sendfile + index 二分 │ KRaft __cluster_metadata │
│ (磁盘 IO) │ + 副本同步 │ + PageCache 读 │ + Topic config 持久化 │
└───────────────┴────────┴────────┴────────┘怎么用这张矩阵:当你遇到一个具体问题,先定位它落在「哪一横、哪一纵」:
- 「我的 Producer 写入卡住了」→ 客户端层 × 生产 → 看
record-error-rate/request-latency-avg,再下钻到协议层(是不是min.insync.replicas不满足)→ 最后到存储层(磁盘是不是满了)。 - 「Consumer 一直 Rebalance」→ 客户端层 × 消费 → 看
rebalance-rate→ 协议层(是不是max.poll.interval.ms超时)→ 存储层一般不背锅。 - 「新增 Topic 报 NOT_CONTROLLER」→ 管理 × 协议层(Controller 切换中)→ 等待几秒重试。
4. 三个隐藏的内部 Topic(必看)
Kafka 有几个带双下划线前缀的内部 Topic,用户不能直接创建 / 删除 / 改配置,但它们承担了 Kafka 自身的「关键基础设施」角色。理解它们 = 理解 Kafka 元数据 / 状态是怎么持久化的。
4.1 __consumer_offsets —— 消费组进度的持久化账本
- 作用:存储所有 Consumer Group 的 Offset 提交记录(
OffsetCommit请求最终落地到这里)和 Group 元数据快照。 - 结构:默认 50 个分区(
offsets.topic.num.partitions,建议保持默认)、副本数offsets.topic.replication.factor(生产建议 3)。 - 路由规则:
Math.abs(group.id.hashCode()) % 50→ 决定 Group 落在哪个分区,该分区的 Leader Broker 即该 Group 的 Coordinator。 - 特性:Cleanup Policy =
compact,因为我们只关心每个(group, topic, partition)的最新一条 Offset,旧的 Compact 掉省空间(详见第 14 章 Log Compaction)。 - 生活类比:邮局后台的「家庭邮箱进度本」 —— 张家上次取信取到第 100 号,李家取到第 80 号,王家取到第 65 号,旧的进度被一行一行覆盖。
- 怎么看:
bash
kafka-console-consumer.sh --bootstrap-server localhost:9092 \
--topic __consumer_offsets --partition 25 --from-beginning \
--formatter "kafka.coordinator.group.GroupMetadataManager\$OffsetsMessageFormatter"输出(节选,已脱敏):
[learn-01-group,learn.01.hello,0]::OffsetAndMetadata(offset=2, leaderEpoch=Optional[0],
metadata=, commitTimestamp=1713344050123, expireTimestamp=None)
[learn-01-group,learn.01.hello,1]::OffsetAndMetadata(offset=2, leaderEpoch=Optional[0], ...)每行就是一次 OffsetCommit,「(group, topic, partition) → 已消费到的 Offset + 元数据」。
4.2 __transaction_state —— 事务状态机的持久化日志
- 作用:存储所有事务 Producer 的状态机变迁(
Empty → Ongoing → PrepareCommit → CompleteCommit / CompleteAbort)。 - 结构:默认 50 个分区(
transaction.state.log.num.partitions),副本数transaction.state.log.replication.factor。 - 路由规则:
Math.abs(transactional.id.hashCode()) % 50→ 决定该事务落在哪个分区,该分区 Leader 即该事务的 Transaction Coordinator。 - 特性:Cleanup Policy =
compact,每个transactional.id只保留最新一行状态。 - 生活类比:「公证处的盖章流水」 —— 第 1234 号事务在 12:00:00 进入 Ongoing、12:00:02 进入 PrepareCommit、12:00:03 进入 CompleteCommit,旧状态被新状态覆盖。
- 怎么看:
bash
kafka-dump-log.sh --files /var/kafka/__transaction_state-13/00000000000000000000.log \
--transaction-log-decoder第 13 章会演示一个完整事务的状态机轨迹,本章只需要知道「事务的 ACID 不是凭空保证的,是写在这个 Topic 里的」。
4.3 __cluster_metadata —— KRaft 时代的「集群中央账本」
- 作用:KRaft 模式专属。存储整个集群的元数据变更事件流:Broker 上下线、Topic 创建删除、Partition Reassignment、ISR 变更、Controller 选举 …… 所有元数据变化都是这个 Topic 上的一条「事件」。
- 结构:只有 1 个分区(元数据只能有一条因果序),由 Controller Quorum(3-5 台 Controller 节点)通过 Raft 协议复制。普通 Broker 不持有它的副本,只是订阅者,把事件回放到本地内存里构建 Metadata Cache。
- 特性:写入路径 = Raft 共识(Leader Append → 多数派 Follower Apply → Commit);读取路径 = 普通 Broker 通过
MetadataFetch增量拉取 + 本地回放。 - 生活类比:「集团总部的公司年鉴」 —— 集团里每发生一次组织架构调整、新部门成立、人事任免,都按时间顺序记一笔;各分公司订阅这本年鉴,照着同步本地的「公司组织图」。
- 怎么看:
bash
kafka-metadata-shell.sh --snapshot \
/var/kafka/__cluster_metadata-0/00000000000000000000.log
>> ls /
brokers topics features configs acls
>> cat /topics/learn.02.cluster
{
"name": "learn.02.cluster",
"topicId": "f8a1...",
"partitions": {
"0": {"leader": 1, "replicas": [1,2,3], "isr": [1,2,3], ...},
...
}
}第 10 章我们会用一整章拆解 KRaft —— Controller Quorum、Raft 选举、Snapshot、ZK 迁移路径。本章只需要把它和 __consumer_offsets / __transaction_state 三个 Topic 的「角色定位」记住即可。
4.4 三个内部 Topic 速查表
| 内部 Topic | 引入版本 | 分区数(默认) | 主要作用 | Cleanup 策略 | 关联 Coordinator |
|---|---|---|---|---|---|
__consumer_offsets | 0.9 | 50 | 消费组 Offset & 元数据 | compact | Group Coordinator |
__transaction_state | 0.11 | 50 | 事务 Producer 状态机 | compact | Transaction Coordinator |
__cluster_metadata | 2.8 (KRaft) | 1 | KRaft 集群级元数据日志 | 自定 (Snapshot + Log) | Controller Quorum |
📌 共同特点:都是 Kafka 自己消费自己的典型范式 —— 用一个普通 Topic 的语义来承载关键基础设施数据。这也是「Kafka 是日志平台,不是消息队列」最有说服力的内部佐证:连它自己的元数据都用日志存。
5. 实验:用 kafka-topics.sh --describe 看真实拓扑
5.1 准备一个三副本的 Topic
bash
# 假设已经按根目录 docker-compose.yml 起了 3 broker 集群
docker exec -it kafka1 bash
/opt/kafka/bin/kafka-topics.sh --bootstrap-server localhost:19092 \
--create --topic learn.02.cluster \
--partitions 6 --replication-factor 3 \
--config min.insync.replicas=2 \
--config retention.ms=86400000输出:
Created topic learn.02.cluster.5.2 看分区分布
bash
/opt/kafka/bin/kafka-topics.sh --bootstrap-server localhost:19092 \
--describe --topic learn.02.cluster真实输出(Broker ID = 1/2/3):
Topic: learn.02.cluster TopicId: f8a1xY3pSi2qZ-mQpQ7cRA PartitionCount: 6 ReplicationFactor: 3 Configs: min.insync.replicas=2,retention.ms=86400000,segment.bytes=1073741824
Topic: learn.02.cluster Partition: 0 Leader: 1 Replicas: 1,2,3 Isr: 1,2,3 Elr: LastKnownElr:
Topic: learn.02.cluster Partition: 1 Leader: 2 Replicas: 2,3,1 Isr: 2,3,1 Elr: LastKnownElr:
Topic: learn.02.cluster Partition: 2 Leader: 3 Replicas: 3,1,2 Isr: 3,1,2 Elr: LastKnownElr:
Topic: learn.02.cluster Partition: 3 Leader: 1 Replicas: 1,3,2 Isr: 1,3,2 Elr: LastKnownElr:
Topic: learn.02.cluster Partition: 4 Leader: 2 Replicas: 2,1,3 Isr: 2,1,3 Elr: LastKnownElr:
Topic: learn.02.cluster Partition: 5 Leader: 3 Replicas: 3,2,1 Isr: 3,2,1 Elr: LastKnownElr:把这段输出和本章 §1 的术语 11 件套对照:
| 术语 | 在输出里的位置 |
|---|---|
| Topic | Topic: learn.02.cluster |
| TopicId | KRaft 给的 UUID f8a1xY3pSi2qZ-mQpQ7cRA(删除重建后是新 ID,防止「同名 Topic 数据错乱」) |
| Partition | 6 个分区 0 ~ 5 |
| Replica | Replicas: 1,2,3 表示这个分区在 Broker 1/2/3 上各有一份 |
| Leader | Leader: 1 表示当前对外读写的副本在 Broker 1 |
| Follower | Replicas 里除 Leader 外的就是 Follower |
| ISR | Isr: 1,2,3 表示三副本目前都「同步在线」 |
| Elr / LastKnownElr | KIP-966 引入的「Eligible Leader Replicas」,4.x 新特性,本章先忽略 |
💡 观察分区分布的规律:6 个分区轮转分布在 3 台 Broker,Leader 平均分布(每台扛 2 个 Leader)。这是 Kafka 自动做的负载均衡(详见第 19 章 Preferred Leader Election)。
5.3 看一眼集群级 Broker 信息
bash
/opt/kafka/bin/kafka-broker-api-versions.sh --bootstrap-server localhost:19092 | head -20输出:
kafka1:19092 (id: 1 rack: rack-a) -> (
Produce(0): 0 to 11 [usable: 11],
Fetch(1): 0 to 17 [usable: 17],
...
)
kafka2:19094 (id: 2 rack: rack-b) -> ( ... )
kafka3:19096 (id: 3 rack: rack-c) -> ( ... )可以看到 3 台 Broker、各自的 ID、机架(rack)和支持的 API 版本范围。这就是 AdminClient 调用 DescribeCluster 拿到的信息的简化版。
5.4 看内部 Topic
加上 --exclude-internal=false(默认就是这样)--list 一下:
bash
/opt/kafka/bin/kafka-topics.sh --bootstrap-server localhost:19092 --list输出:
__consumer_offsets
__transaction_state
learn.01.hello
learn.02.cluster注意 KRaft 模式下 __cluster_metadata 不会出现在普通 Topic 列表(它由 Controller Quorum 单独管理),要看它需要:
bash
/opt/kafka/bin/kafka-metadata-quorum.sh --bootstrap-server localhost:19092 describe --status输出:
ClusterId: MkU3OEVBNTcwNTJENDM2Qk
LeaderId: 3
LeaderEpoch: 12
HighWatermark: 1456
MaxFollowerLag: 0
MaxFollowerLagTimeMs: 0
CurrentVoters: [1, 2, 3]
CurrentObservers: []LeaderId: 3→ Active Controller 是 Broker 3。CurrentVoters: [1,2,3]→ 三台 Broker 都是 Controller Quorum 成员。HighWatermark: 1456→__cluster_metadata的 HW,意思是「截至此元数据事件之前的所有变更都已被多数派确认」。
5.5 配套 Python 探针:把上面所有信息一行命令拉回来
参见本章配套 02_architecture/code/cluster_probe.py:
bash
python 02_architecture/code/cluster_probe.py它会把 Cluster ID、Active Controller、Broker 列表、所有 Topic(含内部)、每个 Topic 的 Partition 详情都打印成漂亮的 ASCII 表格。这是你「接手任何 Kafka 集群的第一件事」。
6. 与其他 MQ 的「角色对照表」
第 1 章 §1.4 已经做过功能对比,这里专门做「架构角色 / 元数据 / 协调者」层面的对应,便于跨技术栈同学迅速建立心智模型。
6.1 角色对应一览
| 角色概念 | Kafka | RabbitMQ | RocketMQ | Pulsar |
|---|---|---|---|---|
| 服务节点 | Broker | Node(Erlang 节点) | Broker | Broker(无状态) |
| 存储单元 | Partition / Segment | Queue | ConsumeQueue + CommitLog | Ledger(BookKeeper) |
| 逻辑容器 | Topic | Exchange + Queue | Topic + Tag | Topic |
| 元数据管理 | KRaft (__cluster_metadata) | Mnesia 内嵌 KV | NameServer | ZooKeeper |
| 消费组协调 | Group Coordinator | 无(Queue 即组) | Broker 内 Coordinator | Broker(订阅) |
| 进度存储 | __consumer_offsets Topic | Queue 内 ack 状态 | OffsetTable 文件 | __consumer_offsets(兼容 Kafka) |
| 事务协调 | Transaction Coordinator + __transaction_state | TxChannel | Half Message + 二次确认 | TxnLog |
| 副本机制 | Leader-Follower(ISR) | Mirrored Queue / Quorum Queue | Master-Slave / DLedger | BookKeeper Ensemble |
| HA 选举 | KRaft Raft | Erlang RPC + 选举算法 | NameServer 心跳 | ZK 临时节点 |
6.2 三个最容易让人懵的对应点
- Kafka Topic ≠ RabbitMQ Topic:
- Kafka Topic 是逻辑容器 + N 个 Partition,路由只能按 Key Hash。
- RabbitMQ「Topic Exchange」只是路由模式之一,按 routing key 通配符匹配,路由能力强大得多。
- Kafka 没有「队列」概念,但 Partition 像极了「分片队列」:
- 想要 RabbitMQ 那种「Queue」语义?把 Topic 设为 1 个 Partition + 1 个 Consumer Group + 1 个 Consumer 就是。但你放弃了所有并行度。
- Kafka 的 Coordinator 在 Broker 上,不像 Pulsar 那样独立部署:
- Pulsar 是「Broker(无状态计算)+ BookKeeper(存储)+ ZK(元数据)」三层独立。
- Kafka 的 Broker 同时承担:存储 + 计算 + Group Coordinator + Tx Coordinator + 可选 Controller,5 合 1。
📌 架构哲学差异:Kafka 选择「Broker 5 合 1 + 极简元数据」换来运维简单、单机吞吐爆表;Pulsar 选择「计算存储分离」换来弹性扩缩、多租户隔离。没有银弹,只有取舍。
7. 本章小结
┌──────────────────────────────────────────────────────────────────┐
│ 第 2 章 核心要点 │
├──────────────────────────────────────────────────────────────────┤
│ │
│ ① 11 个核心术语: │
│ Producer 寄信人 / Broker 分拣中心 / Consumer 收信人 │
│ Topic 信件柜 / Partition 流水线 / Replica 备份流水线 │
│ Offset 流水号 / Consumer Group 家庭邮箱 │
│ Group Coordinator 家庭班长 / Transaction Coordinator 公证处 │
│ Controller 集团调度中心 / KRaft 自带流水的元数据账本 │
│ │
│ ② 一条消息生命周期 12 步: │
│ Producer 攒批 → Leader 追加 → Follower Fetch → ISR LEO 推 HW │
│ → 回 ACK → Consumer Fetch(≤ HW)→ 业务处理 → OffsetCommit │
│ │
│ ③ 三横三纵分层: │
│ 存储层 / 协议层 / 客户端层 │
│ × 生产 / 消费 / 管理 = 9 个能力格子 │
│ │
│ ④ 3 个内部 Topic: │
│ __consumer_offsets : 50 分区, compact, Group Coordinator │
│ __transaction_state : 50 分区, compact, Tx Coordinator │
│ __cluster_metadata : 1 分区, KRaft Raft 复制 │
│ │
│ ⑤ 实验确认: │
│ kafka-topics.sh --describe → 肉眼看 Leader / Replicas / ISR │
│ kafka-metadata-quorum.sh → 肉眼看 Active Controller │
│ cluster_probe.py → 一键拉全部拓扑 │
│ │
│ ⑥ 与其他 MQ 角色对照: │
│ Kafka Broker = 5 合 1(存储 + 计算 + Group Coord + Tx Coord │
│ + 可选 Controller) │
│ Pulsar 选择存算分离,Kafka 选择极简集成 │
│ │
└──────────────────────────────────────────────────────────────────┘8. 面试高频题
Q1:详细说一下 Kafka 的整体架构,以及一条消息从 Producer 到 Consumer 的完整链路?
考察点:能否在 5 分钟内结构化、不遗漏地讲清楚整个集群,而不是只说「Producer 发给 Broker,Consumer 来拉」。
标准答案(按四段式作答):
集群结构(先讲组件):
- 客户端层:Producer / Consumer / AdminClient(基于 Kafka 协议的 TCP 客户端)。
- 服务端集群:3+ 台 Broker(KRaft 模式下其中 3-5 台兼任 Controller,组成 Quorum)。每台 Broker 维护一组 Partition 副本(Leader 或 Follower)。
- 协调者:每个 Broker 都可能扮演某些 Consumer Group 的 Group Coordinator、某些事务的 Transaction Coordinator,由
__consumer_offsets/__transaction_state的分区 Leader 决定。 - 元数据:KRaft 模式下集群元数据存在
__cluster_metadataTopic,由 Controller Quorum Raft 复制。
写入链路(Producer → Leader → Follower → ACK):
- Producer
produce()把消息序列化、Partitioner 选分区、攒到本地 batch(linger.ms/batch.size); - 一批成形后通过 TCP 发
ProduceRequest给该分区 Leader; - Leader 顺序追加到
.log,分配 Offset; - Follower 后台 Fetcher 线程不停
FetchRequestLeader,把数据复制到本地,并把自己的 LEO 汇报给 Leader; - Leader 把 ISR 中所有副本的 LEO 取 min 推进 HW;
- 当
acks=all且 HW ≥ 本批 offset 时,Leader 才回ProduceResponse给 Producer。
- Producer
消费链路(Consumer → Leader → 业务处理 → OffsetCommit):
- Consumer 启动时向任意 Broker
FindCoordinator找到自己 Group 的 Coordinator; JoinGroup+SyncGroup加入组、拿到分区分配;- 循环
poll()→ 实际是向各分区 Leader 发FetchRequest(只能读到 HW 之前的数据); - 业务处理完调用
commit()→OffsetCommit写到__consumer_offsets; - 周期性
Heartbeat保活,否则 Coordinator 触发 Rebalance 把分区分给别人。
- Consumer 启动时向任意 Broker
元数据 / 故障切换:
- Controller 监控所有 Broker 的注册(KRaft 下通过
BrokerRegistrationHeartbeat); - Broker 挂了,Controller 把它身上的所有 Leader 切换到 ISR 中的 Follower,更新
__cluster_metadata; - 普通 Broker 通过
MetadataFetch拉新元数据,通知客户端切流量。
- Controller 监控所有 Broker 的注册(KRaft 下通过
加分项:
- 能画出本章 §2.3 那张时序图,并能说出第 7 步「HW = min(ISR.LEO)」的细节。
- 能讲清「写入路径要等多副本,读路径只读 Leader」的非对称性。
- 能补一句「
acks=0/1/all决定 Producer 等待程度」「min.insync.replicas决定可用副本最低门槛」。
与其他 MQ 对比:
- RabbitMQ 是 Push 模型 + Queue 节点持久化,Broker 必须维护「消息发给谁了 / 是否 ACK」状态。
- RocketMQ 写入 CommitLog 是单文件追加(vs Kafka 的多 Partition 多文件),消费时通过 ConsumeQueue 索引定位。
- Pulsar 的 Broker 是无状态计算节点,所有数据写到 BookKeeper Bookie,Broker 故障切换不涉及数据迁移。
Q2:Kafka 的 Controller 是干什么的?KRaft 之前和之后有什么不同?
考察点:是否真的理解 Kafka「集群级元数据管理」的来龙去脉,而不只是会背「Controller 管理元数据」。
标准答案:
Controller 的职责(无论 ZK 还是 KRaft 时代都不变):
- 监听 Broker 上线 / 下线(注册);
- 触发 Partition Leader 选举(Broker 挂了、Reassignment 完成、
unclean.leader.election等); - 维护 ISR 变更(Follower 掉队 / 追上);
- 处理 Topic 创建 / 删除 / Reassignment 等管理请求;
- 把元数据变化广播给所有 Broker。
ZK 时代(Kafka < 2.8):
- 整个集群里所有 Broker 都是 Controller 候选,启动时抢 ZK 的
/controller临时节点,谁抢到谁就是当前 Controller。 - 元数据存储在 ZK 上的多个路径(
/brokers/ids、/brokers/topics、/admin/reassign_partitions等)。 - 缺点:① 元数据上限受 ZK 限制(生产经验上 ~20 万分区就到瓶颈);② Controller 切换时要从 ZK 重新拉全量元数据,秒级不可用;③ 多个组件(Broker / KafkaUI / Kafka-Connect)都要装 ZK 客户端,运维复杂;④ 元数据变更不是事件溯源,难以做完整 audit。
- 整个集群里所有 Broker 都是 Controller 候选,启动时抢 ZK 的
KRaft 时代(Kafka ≥ 2.8 预览,≥ 3.3 GA,4.0 全面替代):
- 引入专门的 Controller 节点角色(可以与 Broker 合并部署,也可以独立),3-5 台组成一个 Raft Quorum,由 Raft 选出 Active Controller,其它是 Standby。
- 元数据存到一个特殊内部 Topic
__cluster_metadata,每条变更是一条事件,Raft 复制到所有 Voter。 - 普通 Broker 不再连 ZK,而是通过
MetadataFetch增量拉__cluster_metadata上的事件,回放到本地构建元数据缓存。 - 优势:① 元数据上限 → 百万分区级;② Controller 切换秒级(Raft 任期切换 + Standby 已经有完整状态);③ 去 ZK,运维只剩 Kafka 一套;④ 元数据天然事件溯源,便于审计 / 快照 / 重放。
迁移路径:
- 3.x 提供「双模式」 —— 你可以选 ZK 或 KRaft 启动;
- 3.4+ 提供「Migration Tool」:先把 Broker 升级到双模式,再把 ZK 元数据搬到 KRaft,最后下线 ZK;
- 4.0 起 KRaft 是唯一选择,老集群必须先迁移。
加分项:
- 能讲「Controller Quorum 与 Broker 角色的关系」:通过
process.roles配置,可以是broker/controller/broker,controller(合并)。 - 能讲「Active Controller 怎么看」:
kafka-metadata-quorum.sh describe --status看LeaderId。 - 能讲「为什么是 Raft 而不是 Paxos / ZAB」:Raft 算法工程实现简单、有现成参考实现、Kafka 团队熟悉日志语义。
与其他元数据方案对比:
- ZooKeeper:通用 KV,强一致 ZAB,但元数据放大、客户端连接重;
- etcd:通用 KV,Raft,K8s 生态用得最多;
- Pulsar:用 ZK 管元数据 + BookKeeper 管数据;
- Kafka KRaft:把「元数据」做成自己的 Topic,用自己消费自己实现自洽。
Q3:什么是 ISR / HW / LEO?它们的关系是怎样的?
考察点:副本机制的三个核心数字关系,是后面理解 acks=all、unclean.leader.election、Leader Epoch 的基础。
标准答案:
三个名词定义:
- LEO(Log End Offset):单个副本本地日志写到的「下一条要写入」的 Offset。每个副本(含 Leader 和 Follower)都有自己的 LEO。
- HW(High Watermark):Partition 当前ISR 中所有副本都已经追上的最高 Offset,是 Consumer 可见的「水位线」 —— 超过 HW 的消息 Consumer 看不见。
- ISR(In-Sync Replicas):当前与 Leader 保持同步的副本集合(含 Leader 自己)。判定标准:
replica.lag.time.max.ms(默认 30s)内拉过最新 offset。
三者的关系(图示):
Leader 日志: [0][1][2][3][4][5][6] LEO=7 Foll-A 日志: [0][1][2][3][4][5] LEO=6 (在 ISR) Foll-B 日志: [0][1][2][3][4] LEO=5 (在 ISR) Foll-C 日志: [0][1] LEO=2 (掉队 30s, 已踢出 ISR) HW = min(LEO of ISR) = min(7, 6, 5) = 5 Consumer 可见: offset 0~4 (即 < HW)- HW = ISR 中所有副本 LEO 的最小值。
unclean.leader.election=false(默认)下,Leader 挂了只能从 ISR 选新 Leader → HW 不会回退(数据安全但可能不可用);unclean.leader.election=true时可以从 ISR 外(即 OSR)选 Leader → 已 ack 消息可能丢失(数据可能不安全但提高可用性)。
生活类比:
- LEO = 每条快递员(副本)「自己手上拿到的最新一封」编号;
- HW = 「所有可信快递员都拿到了的最高编号」 —— 只有低于这个编号的信,邮局才允许收件人看(防止你看到的信其实有些快递员没拿到,万一 Leader 挂了不就丢了);
- ISR = 「可信快递员名单」 —— 太久没拉信件的快递员(掉队 30s 以上)会被剔除,等他追上再重新加入。
典型故障场景:
- 写入 acks=all → 必须 HW 推进到本次写入 offset 才回 ACK → 保证「ACK 过的消息至少在 ISR 所有副本都有」。
- Leader 挂了 → Controller 从 ISR 中选新 Leader(默认是 Replicas 列表里的第一个仍在 ISR 中的)→ 新 Leader 把 HW 之后未确认的部分截断(详见第 9 章 Leader Epoch 解决「截断丢数据」)。
加分项:
- 能补一句「
min.insync.replicas配合acks=all才是真不丢消息组合」 —— 例如 RF=3 + min.isr=2 + acks=all,最多容忍 1 台挂掉。 - 能讲
LogStartOffset(LSO)和 HW 的关系:LSO 是「最早还能读到的 offset」(被删除策略清理过),HW 是「最新能读到的 offset + 1」。 - 能讲
LeaderEpoch如何防止「截断 + 旧 Leader 复活 → 数据分叉」(第 9 章)。
与其他 MQ 对比:
- RocketMQ 主从同步有「同步双写」和「异步复制」两种模式,没有 ISR 这种「动态可信副本集合」;
- Pulsar 的副本(Bookie)以 Quorum 写入(Ensemble + Write Quorum + Ack Quorum),更细粒度;
- RabbitMQ Mirrored Queue 是镜像复制,新版 Quorum Queue 用 Raft 类似机制。
Q4:Group Coordinator 是怎么定位的?为什么 __consumer_offsets 默认要 50 个分区?
考察点:消费组协议的协调者定位机制,以及为什么 Kafka 选择「内部 Topic 哈希分片」这种设计。
标准答案:
Group Coordinator 定位算法:
- Consumer 启动时向
bootstrap.servers中的任意 Broker 发FindCoordinatorRequest,传group.id。 - 任意 Broker 都能算出公式
partitionId = Math.abs(group.id.hashCode()) % offsets.topic.num.partitions(默认% 50)。 - 这个分区在
__consumer_offsets里的 Leader 所在 Broker,就是该 Group 的 Coordinator。 - Broker 把 Coordinator 的 host/port 返回给 Consumer,后者直接连过去发
JoinGroup。
- Consumer 启动时向
为什么是 50 个分区:
- 均匀分散负载:50 个分区均匀分布在 N 台 Broker 上,避免某一台扛所有 Group 的协调流量。3 台 Broker 集群下每台扛 17 个 Coordinator 角色。
- 限制单分区上限:50 个 Group Coordinator 分区 + 每个分区 RF=3,总共 150 个副本,对元数据来说可控。
- 避免数据热点:每个 Group 的 OffsetCommit 流量集中在自己的协调分区,分区多了之后单个分区的写入压力可控。
- 历史选择:50 是经验拍出来的,社区一般不建议改。如果集群超大(万级 Group),可以适度提高到 100~200。
协调者切换:
- Coordinator 所在 Broker 挂了 → 该
__consumer_offsets分区会触发 Leader 选举 → 新 Leader 所在 Broker 自动接管 Coordinator 职责。 - Consumer 收到
NOT_COORDINATOR错误 → 重新FindCoordinator→ 切到新 Coordinator。 - 这就是 Kafka 不需要专门的「Coordinator 选举服务」的原因 —— 所有协调状态都搭车在普通 Topic 的副本机制上。
- Coordinator 所在 Broker 挂了 → 该
生活类比:
- 集团里 50 个「家庭服务部门」,按家庭姓氏首字母分;张家归 Z 部门管,李家归 L 部门管 …… 部门负责人换人不影响家庭对应关系;
- 任何一个集团接待员都能告诉你「张家归哪个部门」,因为算法(哈希)是公开的。
加分项:
- 能讲清「Coordinator 切换会触发该组所有成员重新 JoinGroup」 —— Heartbeat 失败 → 重连 → 重新加入组(一次轻量 Rebalance)。
- 能补「事务用同样的设计」:
__transaction_state也是 50 分区,根据transactional.id哈希定位 Transaction Coordinator。
与其他 MQ 对比:
- RabbitMQ:Queue 本身就是协调单位,没有「组协调者」一说,靠 Erlang 节点间复制;
- RocketMQ:消费进度由 Broker 直接维护在 OffsetTable 文件,Consumer Pull 时查询,不需要专门 Coordinator;
- Pulsar:Subscription 状态由 Broker 维护(前端 Broker 负载均衡),最终一致性写到 BookKeeper。
Q5:什么是 __consumer_offsets / __transaction_state / __cluster_metadata?它们是普通 Topic 吗?
考察点:能否区分 Kafka 的 3 个内部 Topic、它们的用途、与普通 Topic 的差异。
标准答案:
| 维度 | __consumer_offsets | __transaction_state | __cluster_metadata |
|---|---|---|---|
| 引入版本 | 0.9 | 0.11 | 2.8 (KRaft) |
| 默认分区数 | 50 | 50 | 1(不可改) |
| 副本数 | offsets.topic.replication.factor(推荐 3) | transaction.state.log.replication.factor(推荐 3) | Controller Quorum 决定(3-5) |
| 用途 | 存 Consumer Group 的 OffsetCommit + Group 元数据 | 存事务 Producer 的状态机变迁 | 存集群元数据事件流(Topic / Partition / Broker / ACL …) |
| 清理策略 | compact | compact | Snapshot + Log(KRaft 自管) |
| 谁在写 | Group Coordinator(写 Consumer 提交的 Offset) | Transaction Coordinator | Active Controller |
| 用户能直接写吗 | 不能(受保护) | 不能 | 不能(普通 Broker 都不持有副本,只是订阅) |
它们和普通 Topic 的相似 vs 差异:
- 相似处:底层存储格式、副本机制、Compaction 策略,都是普通 Kafka Topic 那一套 —— 这就是「Kafka 自己消费自己」的体现,证明日志抽象的强大。
- 差异处:
- 用户不能通过
kafka-topics.sh --create / --delete操作(只能改部分配置); - 默认在
--list输出里被标为「internal」(加--exclude-internal隐藏); - 写入入口是协议层定义的特定 Request(OffsetCommit / EndTxn / RegisterBroker)而非普通 ProduceRequest。
- 用户不能通过
实战要点:
- 生产环境必须把
offsets.topic.replication.factor设为 3(默认 1,单 Broker 集群上的「假默认值」)—— 否则 Coordinator Broker 一挂,所有 Group 的 Offset 进度全丢。 __consumer_offsets的 Compaction 是必须的,否则随便一个频繁 Commit 的应用都能把分区撑爆。- KRaft 模式下,
__cluster_metadata是集群的「单点故障防护」 —— 必须保证 Controller Quorum 多数派存活,少了 (N+1)/2 就挂了。
加分项:
- 能讲
__consumer_offsets的 Key 编码:(version, group.id, topic, partition),Value 编码:(offset, leaderEpoch, metadata, commitTimestamp, expireTimestamp)。 - 能用
kafka-console-consumer.sh --formatter "kafka.coordinator.group.GroupMetadataManager$OffsetsMessageFormatter"直接把__consumer_offsetsdump 出来。 - 能讲 KRaft 的 Snapshot 机制:日志太长时,Active Controller 把当前元数据状态序列化成 Snapshot 文件,新 Controller 启动时先加载 Snapshot 再回放后续日志。
与其他 MQ 对比:
- RabbitMQ:元数据用 Mnesia(Erlang 内置 KV),Offset 状态在 Queue 内部数据结构;
- RocketMQ:元数据放 NameServer 内存 + 定时持久化,Offset 放 Broker 本地 OffsetTable JSON 文件;
- Pulsar:元数据 + Offset 都放 ZK,数据放 BookKeeper。
Q6(加分题):Kafka 客户端是怎么知道某个 Partition 的 Leader 在哪台 Broker 上的?
考察点:Kafka 协议里 MetadataRequest 的角色,理解客户端 / 集群之间的协议层互动。
标准答案:
- 首次启动:Producer / Consumer 用
bootstrap.servers里随便找一台 Broker 连上去,发MetadataRequest(topics)。 - Broker 回包:返回该 Topic 所有 Partition 的元数据,包括每个 Partition 的
Leader、Replicas、ISR、以及对应 Broker 的host、port。 - 客户端建立 TCP 直连:客户端按 Partition Leader 的位置,和需要的 Broker 建立 TCP 连接(Producer 发到 Leader,Consumer 拉自 Leader)。
- 元数据失效再次刷新:触发条件包括:
metadata.max.age.ms到期(默认 5 分钟)周期性刷新;- 收到
NOT_LEADER_FOR_PARTITION/LEADER_NOT_AVAILABLE错误立刻重新MetadataRequest; - 收到新元数据版本(Broker 在响应里带
controller_epoch/metadata_version);
- 元数据传播:Active Controller 把元数据变更写到
__cluster_metadata,普通 Broker 通过MetadataFetch拉取并更新自己的 Metadata Cache,客户端从任意 Broker 拉到的元数据都最终一致。
加分项:
- 能讲「为什么
bootstrap.servers不需要写完整集群」 —— 只要能联通一个活的 Broker,客户端就能拿到完整集群元数据;推荐写 2-3 个做容灾。 - 能讲「集群滚动升级时元数据 staleness」 —— 客户端可能短暂连到不再是 Leader 的 Broker,靠 NOT_LEADER 错误自愈。
- 能讲 KIP-700 / KIP-714「Client Telemetry」让 Broker 拿到客户端版本与指标。
📌 下一章预告:第 3 章我们会回到「怎么用」的视角,拆解
kafka-*.sh命令行家族 + Pythonconfluent-kafka的三大 API(Producer / Consumer / AdminClient),让你能在自己的集群上把第 1、2 章学到的术语亲手敲一遍。
🎬 可视化演示
演示加载缓慢或样式异常?点此在新标签页打开 ↗
💻 示例代码
python
"""
第 2 章 · Cluster Probe —— 把整个 Kafka 集群拓扑「一键打印」
用法:
python 02_architecture/code/cluster_probe.py
python 02_architecture/code/cluster_probe.py --bootstrap localhost:9092
python 02_architecture/code/cluster_probe.py --bootstrap localhost:9092 --include-internal
它会做四件事:
1. 连上 Kafka 集群,拿 Cluster ID + Active Controller ID。
2. 列出所有 Broker(id / host / port / rack)。
3. 列出所有 Topic(默认隐藏内部 Topic,可用 --include-internal 显示
__consumer_offsets / __transaction_state / __cluster_metadata 等)。
4. 对每个 Topic 打印所有 Partition:partition_id / leader / replicas / isr。
输出全部用 `tabulate` 渲染成 ASCII 表格,方便复制贴到工单 / 文档里。
连不上 Broker 时会给出友好提示(端口、防火墙、bootstrap 地址等常见排查方向)。
依赖:
pip install -r requirements.txt # 至少需要 confluent-kafka>=2.5 + tabulate>=0.9
预期输出(节选,3 broker × KRaft 集群,已建若干 demo Topic):
┌────────────────────────── 集群概览 ──────────────────────────┐
│ Cluster ID : MkU3OEVBNTcwNTJENDM2Qk │
│ Bootstrap : localhost:9092 │
│ Total Brokers : 3 │
│ Active Controller : Broker 3 (KRaft Active) │
│ Total Topics : 4 (含内部 1) │
│ Total Partitions : 56 │
└──────────────────────────────────────────────────────────────┘
── Brokers ───────────────────────────────────
ID Host Port Rack
---- ----------- ------ ------
1 kafka1 19092 rack-a
2 kafka2 19094 rack-b
3 kafka3 19096 rack-c
── Topic: learn.02.cluster (RF=3, Partitions=6, Internal=No) ──
Partition Leader Replicas ISR OOS-Replicas
----------- ------ -------- ------ --------------
0 1 1,2,3 1,2,3 ()
1 2 2,3,1 2,3,1 ()
...
"""
from __future__ import annotations
import argparse
import sys
import time
from typing import List, Tuple
try:
from confluent_kafka.admin import AdminClient, ConfigResource
from confluent_kafka import KafkaException
except ImportError as e:
print(f"[FATAL] 缺少依赖:{e}\n请先安装:pip install -r requirements.txt")
sys.exit(1)
try:
from tabulate import tabulate
except ImportError:
print("[FATAL] 缺少依赖 tabulate,请先 pip install tabulate>=0.9")
sys.exit(1)
INTERNAL_TOPIC_PREFIXES = ("__",) # __consumer_offsets / __transaction_state ...
def parse_args() -> argparse.Namespace:
p = argparse.ArgumentParser(
description="探针:打印 Kafka 集群拓扑(Broker + Topic + Partition)。"
)
p.add_argument(
"--bootstrap",
default="localhost:9092",
help="Kafka bootstrap.servers (默认 localhost:9092)",
)
p.add_argument(
"--include-internal",
action="store_true",
help="同时显示内部 Topic(__consumer_offsets 等)",
)
p.add_argument(
"--timeout",
type=float,
default=10.0,
help="集群元数据获取超时秒数 (默认 10s)",
)
return p.parse_args()
def is_internal(topic: str) -> bool:
return any(topic.startswith(pfx) for pfx in INTERNAL_TOPIC_PREFIXES)
def fetch_metadata(admin: AdminClient, timeout: float):
"""连不上时,把常见原因贴出来,避免新手陷入「卡住不动」迷局。"""
try:
md = admin.list_topics(timeout=timeout)
except KafkaException as e:
print(f"\n[ERROR] 无法连接 Kafka 集群:{e}")
print("\n常见排查方向:")
print(" 1. Broker 是否真的在跑? -> docker ps | grep kafka")
print(" 或在 broker 节点上: netstat -lntp | grep 9092")
print(" 2. bootstrap.servers 写对了吗?默认是 localhost:9092;docker compose 时可能是 19092 / 19094 / 19096。")
print(" 3. 防火墙 / 安全组放通 9092 / 9093 端口。")
print(" 4. KRaft 集群刚启动需要 30~60s 才进入 RUNNING 状态,可稍等再试。")
print(" 5. listeners / advertised.listeners 里写的 hostname 客户端能解析吗?")
sys.exit(2)
except Exception as e:
print(f"\n[ERROR] 元数据请求失败:{e}")
sys.exit(2)
return md
def render_overview(md, bootstrap: str, include_internal: bool) -> None:
cluster_id = md.cluster_id or "(unknown, KRaft 早期版本可能不返回)"
controller = md.controller_id # 这是当前 metadata response 的 controller id
total_topics = len(md.topics)
visible_topics = [t for t in md.topics if include_internal or not is_internal(t)]
total_parts = sum(len(t.partitions) for t in md.topics.values())
lines = [
("Cluster ID", cluster_id),
("Bootstrap", bootstrap),
("Total Brokers", len(md.brokers)),
("Active Controller", f"Broker {controller}" if controller >= 0 else "(unknown)"),
(
"Total Topics",
f"{total_topics} (显示 {len(visible_topics)} | 隐藏内部 {total_topics - len(visible_topics)})",
),
("Total Partitions", total_parts),
]
print()
print("┌" + "─" * 26 + " 集 群 概 览 " + "─" * 26 + "┐")
for k, v in lines:
line = f"│ {k:<19}: {v}"
print(line + " " * max(0, 65 - len(line)) + "│")
print("└" + "─" * 65 + "┘")
def render_brokers(md) -> None:
print("\n── Brokers ──────────────────────────────────────")
rows = []
for bid, b in sorted(md.brokers.items()):
rack = getattr(b, "rack", None) or "-"
rows.append([bid, b.host, b.port, rack])
print(tabulate(
rows,
headers=["ID", "Host", "Port", "Rack"],
tablefmt="simple",
))
def render_topics(md, include_internal: bool) -> None:
topic_names = sorted(md.topics.keys())
if not include_internal:
topic_names = [t for t in topic_names if not is_internal(t)]
if not topic_names:
print("\n[INFO] 当前集群没有任何用户 Topic。试试:")
print(" kafka-topics.sh --bootstrap-server <bs> --create --topic learn.02.demo --partitions 6 --replication-factor 3")
return
for tname in topic_names:
t = md.topics[tname]
if t.error is not None:
print(f"\n── Topic: {tname} ── ⚠️ 元数据错误:{t.error}")
continue
partitions = sorted(t.partitions.values(), key=lambda x: x.id)
rf = len(partitions[0].replicas) if partitions else 0
rows = []
for p in partitions:
replicas = sorted(p.replicas)
isr = sorted(p.isrs)
oos = [r for r in replicas if r not in isr] # Out-of-Sync Replicas
rows.append([
p.id,
p.leader if p.leader >= 0 else "(none)",
",".join(map(str, replicas)),
",".join(map(str, isr)) if isr else "(empty)",
",".join(map(str, oos)) if oos else "()",
])
flag_internal = "Yes" if is_internal(tname) else "No"
print(f"\n── Topic: {tname} (RF={rf}, Partitions={len(partitions)}, Internal={flag_internal}) ──")
print(tabulate(
rows,
headers=["Partition", "Leader", "Replicas", "ISR", "OOS-Replicas"],
tablefmt="simple",
))
unhealthy = [r for r in rows if r[4] != "()"]
if unhealthy:
print(f" ⚠️ 有 {len(unhealthy)} 个分区存在 Out-of-Sync 副本,建议进一步排查(磁盘 / 网络 / GC)。")
no_leader = [r for r in rows if r[1] == "(none)"]
if no_leader:
print(f" ❌ 有 {len(no_leader)} 个分区当前无 Leader,分区不可读写!")
def main() -> int:
args = parse_args()
print(f"[INFO] 正在连接 Kafka:{args.bootstrap} ...")
admin = AdminClient({
"bootstrap.servers": args.bootstrap,
"client.id": "cluster-probe",
"socket.timeout.ms": int(args.timeout * 1000),
})
t0 = time.time()
md = fetch_metadata(admin, args.timeout)
elapsed = (time.time() - t0) * 1000
print(f"[INFO] 元数据获取成功,耗时 {elapsed:.1f} ms。\n")
render_overview(md, args.bootstrap, args.include_internal)
render_brokers(md)
render_topics(md, args.include_internal)
print()
print("提示:")
print(" • 加 --include-internal 可看 __consumer_offsets / __transaction_state(KRaft 模式下 __cluster_metadata 不在此处显示,需用 kafka-metadata-quorum.sh)。")
print(" • 上面的输出 ≈ 多次 kafka-topics.sh --describe + kafka-broker-api-versions.sh 的合并版本。")
print(" • 把这段输出贴到工单 / 故障复盘文档里,就是一份合格的『集群快照』。\n")
return 0
if __name__ == "__main__":
try:
sys.exit(main())
except KeyboardInterrupt:
print("\n[INFO] 用户取消。")
sys.exit(130)