Skip to content

第 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 共识的元数据库  ║
            ║                                                  ║
            ╚════════════════════════════════════════════════╝

读完本章你应当能:

  1. 复述 11 个核心术语,且每个都讲得出「英文 → 中文 → 类比 → 在集群里的物理 / 逻辑位置」。
  2. 看懂一张大图:从 Producer produce() 到 Consumer poll() 中间发生了什么。
  3. 说出 Kafka 的「三横三纵」:存储层 / 协议层 / 客户端层 × 生产 / 消费 / 管理。
  4. 认识 3 个隐藏内部 Topic__consumer_offsets / __transaction_state / __cluster_metadata,并知道它们各自承载什么职责。
  5. 动手跑一遍 kafka-topics.sh --describe + Python cluster_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)端口。
  • 关键属性:每个 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_offsetsauto.offset.resetseek、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。
  • 核心规则
    1. 同一个 Group 内:每个 Partition 只能被 1 个 Consumer 实例消费 —— 单播 / 负载均衡。
    2. 不同 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 ...
  • 关联术语: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 步完整解释:

动作关键参数 / 配置
1Producer 攒一个 batch 后通过 TCP 发 ProduceRequest 给 Leaderlinger.ms / batch.size / compression.type
2Leader 把 batch 顺序追加到本地 .log,给每条 Record 分配递增 Offsetlog.dirs / segment.bytes
3-6Follower 后台 Fetcher 线程不停拉 Leader 数据,同步到本地replica.fetch.max.bytes / num.replica.fetchers
7每次 Follower Fetch 都汇报自己的 LEO;Leader 把 ISR 中所有副本的 LEO 取 min,推进 HWreplica.lag.time.max.ms
8acks=all 时等 HW ≥ 本次写入 offset 才回 ACK;acks=1 只等 Leader;acks=0 不等acks / min.insync.replicas
9-10Consumer pull 模型,每次 poll() 实际是 FetchRequest;只能读到 HW 之前的数据fetch.min.bytes / fetch.max.wait.ms
11业务代码处理消息(反序列化、入库、调下游 API 等)——
12-13Consumer 把「我已经处理到 offset N+1」写到 __consumer_offsetsenable.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=0Fetch=1Metadata=3OffsetCommit=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_offsets0.950消费组 Offset & 元数据compactGroup Coordinator
__transaction_state0.1150事务 Producer 状态机compactTransaction Coordinator
__cluster_metadata2.8 (KRaft)1KRaft 集群级元数据日志自定 (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 件套对照:

术语在输出里的位置
TopicTopic: learn.02.cluster
TopicIdKRaft 给的 UUID f8a1xY3pSi2qZ-mQpQ7cRA(删除重建后是新 ID,防止「同名 Topic 数据错乱」)
Partition6 个分区 0 ~ 5
ReplicaReplicas: 1,2,3 表示这个分区在 Broker 1/2/3 上各有一份
LeaderLeader: 1 表示当前对外读写的副本在 Broker 1
FollowerReplicas 里除 Leader 外的就是 Follower
ISRIsr: 1,2,3 表示三副本目前都「同步在线」
Elr / LastKnownElrKIP-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 角色对应一览

角色概念KafkaRabbitMQRocketMQPulsar
服务节点BrokerNode(Erlang 节点)BrokerBroker(无状态)
存储单元Partition / SegmentQueueConsumeQueue + CommitLogLedger(BookKeeper)
逻辑容器TopicExchange + QueueTopic + TagTopic
元数据管理KRaft (__cluster_metadata)Mnesia 内嵌 KVNameServerZooKeeper
消费组协调Group Coordinator无(Queue 即组)Broker 内 CoordinatorBroker(订阅)
进度存储__consumer_offsets TopicQueue 内 ack 状态OffsetTable 文件__consumer_offsets(兼容 Kafka)
事务协调Transaction Coordinator + __transaction_stateTxChannelHalf Message + 二次确认TxnLog
副本机制Leader-Follower(ISR)Mirrored Queue / Quorum QueueMaster-Slave / DLedgerBookKeeper Ensemble
HA 选举KRaft RaftErlang RPC + 选举算法NameServer 心跳ZK 临时节点

6.2 三个最容易让人懵的对应点

  1. Kafka Topic ≠ RabbitMQ Topic
    • Kafka Topic 是逻辑容器 + N 个 Partition,路由只能按 Key Hash。
    • RabbitMQ「Topic Exchange」只是路由模式之一,按 routing key 通配符匹配,路由能力强大得多。
  2. Kafka 没有「队列」概念,但 Partition 像极了「分片队列」
    • 想要 RabbitMQ 那种「Queue」语义?把 Topic 设为 1 个 Partition + 1 个 Consumer Group + 1 个 Consumer 就是。但你放弃了所有并行度
  3. 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 来拉」。

标准答案(按四段式作答):

  1. 集群结构(先讲组件)

    • 客户端层: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_metadata Topic,由 Controller Quorum Raft 复制。
  2. 写入链路(Producer → Leader → Follower → ACK)

    • Producer produce() 把消息序列化、Partitioner 选分区、攒到本地 batch(linger.ms / batch.size);
    • 一批成形后通过 TCP 发 ProduceRequest 给该分区 Leader;
    • Leader 顺序追加到 .log,分配 Offset;
    • Follower 后台 Fetcher 线程不停 FetchRequest Leader,把数据复制到本地,并把自己的 LEO 汇报给 Leader;
    • Leader 把 ISR 中所有副本的 LEO 取 min 推进 HW;
    • acks=all 且 HW ≥ 本批 offset 时,Leader 才回 ProduceResponse 给 Producer。
  3. 消费链路(Consumer → Leader → 业务处理 → OffsetCommit)

    • Consumer 启动时向任意 Broker FindCoordinator 找到自己 Group 的 Coordinator;
    • JoinGroup + SyncGroup 加入组、拿到分区分配;
    • 循环 poll() → 实际是向各分区 Leader 发 FetchRequest(只能读到 HW 之前的数据);
    • 业务处理完调用 commit()OffsetCommit 写到 __consumer_offsets
    • 周期性 Heartbeat 保活,否则 Coordinator 触发 Rebalance 把分区分给别人。
  4. 元数据 / 故障切换

    • Controller 监控所有 Broker 的注册(KRaft 下通过 BrokerRegistration Heartbeat);
    • Broker 挂了,Controller 把它身上的所有 Leader 切换到 ISR 中的 Follower,更新 __cluster_metadata
    • 普通 Broker 通过 MetadataFetch 拉新元数据,通知客户端切流量。

加分项

  • 能画出本章 §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 管理元数据」。

标准答案

  1. Controller 的职责(无论 ZK 还是 KRaft 时代都不变):

    • 监听 Broker 上线 / 下线(注册);
    • 触发 Partition Leader 选举(Broker 挂了、Reassignment 完成、unclean.leader.election 等);
    • 维护 ISR 变更(Follower 掉队 / 追上);
    • 处理 Topic 创建 / 删除 / Reassignment 等管理请求;
    • 把元数据变化广播给所有 Broker。
  2. 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。
  3. 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 一套;④ 元数据天然事件溯源,便于审计 / 快照 / 重放。
  4. 迁移路径

    • 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 --statusLeaderId
  • 能讲「为什么是 Raft 而不是 Paxos / ZAB」:Raft 算法工程实现简单、有现成参考实现、Kafka 团队熟悉日志语义。

与其他元数据方案对比

  • ZooKeeper:通用 KV,强一致 ZAB,但元数据放大、客户端连接重;
  • etcd:通用 KV,Raft,K8s 生态用得最多;
  • Pulsar:用 ZK 管元数据 + BookKeeper 管数据;
  • Kafka KRaft:把「元数据」做成自己的 Topic,用自己消费自己实现自洽。

Q3:什么是 ISR / HW / LEO?它们的关系是怎样的?

考察点:副本机制的三个核心数字关系,是后面理解 acks=allunclean.leader.election、Leader Epoch 的基础。

标准答案

  1. 三个名词定义

    • 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。
  2. 三者的关系(图示)

    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 消息可能丢失(数据可能不安全但提高可用性)。
  3. 生活类比

    • LEO = 每条快递员(副本)「自己手上拿到的最新一封」编号;
    • HW = 「所有可信快递员都拿到了的最高编号」 —— 只有低于这个编号的信,邮局才允许收件人看(防止你看到的信其实有些快递员没拿到,万一 Leader 挂了不就丢了);
    • ISR = 「可信快递员名单」 —— 太久没拉信件的快递员(掉队 30s 以上)会被剔除,等他追上再重新加入。
  4. 典型故障场景

    • 写入 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 哈希分片」这种设计。

标准答案

  1. 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
  2. 为什么是 50 个分区

    • 均匀分散负载:50 个分区均匀分布在 N 台 Broker 上,避免某一台扛所有 Group 的协调流量。3 台 Broker 集群下每台扛 17 个 Coordinator 角色。
    • 限制单分区上限:50 个 Group Coordinator 分区 + 每个分区 RF=3,总共 150 个副本,对元数据来说可控。
    • 避免数据热点:每个 Group 的 OffsetCommit 流量集中在自己的协调分区,分区多了之后单个分区的写入压力可控。
    • 历史选择:50 是经验拍出来的,社区一般不建议改。如果集群超大(万级 Group),可以适度提高到 100~200。
  3. 协调者切换

    • Coordinator 所在 Broker 挂了 → 该 __consumer_offsets 分区会触发 Leader 选举 → 新 Leader 所在 Broker 自动接管 Coordinator 职责。
    • Consumer 收到 NOT_COORDINATOR 错误 → 重新 FindCoordinator → 切到新 Coordinator。
    • 这就是 Kafka 不需要专门的「Coordinator 选举服务」的原因 —— 所有协调状态都搭车在普通 Topic 的副本机制上
  4. 生活类比

    • 集团里 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.90.112.8 (KRaft)
默认分区数50501(不可改)
副本数offsets.topic.replication.factor(推荐 3)transaction.state.log.replication.factor(推荐 3)Controller Quorum 决定(3-5)
用途存 Consumer Group 的 OffsetCommit + Group 元数据存事务 Producer 的状态机变迁存集群元数据事件流(Topic / Partition / Broker / ACL …)
清理策略compactcompactSnapshot + Log(KRaft 自管)
谁在写Group Coordinator(写 Consumer 提交的 Offset)Transaction CoordinatorActive Controller
用户能直接写吗不能(受保护)不能不能(普通 Broker 都不持有副本,只是订阅)

它们和普通 Topic 的相似 vs 差异

  • 相似处:底层存储格式、副本机制、Compaction 策略,都是普通 Kafka Topic 那一套 —— 这就是「Kafka 自己消费自己」的体现,证明日志抽象的强大。
  • 差异处
    • 用户不能通过 kafka-topics.sh --create / --delete 操作(只能改部分配置);
    • 默认在 --list 输出里被标为「internal」(加 --exclude-internal 隐藏);
    • 写入入口是协议层定义的特定 Request(OffsetCommit / EndTxn / RegisterBroker)而非普通 ProduceRequest。

实战要点

  1. 生产环境必须把 offsets.topic.replication.factor 设为 3(默认 1,单 Broker 集群上的「假默认值」)—— 否则 Coordinator Broker 一挂,所有 Group 的 Offset 进度全丢。
  2. __consumer_offsets 的 Compaction 是必须的,否则随便一个频繁 Commit 的应用都能把分区撑爆。
  3. 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_offsets dump 出来。
  • 能讲 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 的角色,理解客户端 / 集群之间的协议层互动。

标准答案

  1. 首次启动:Producer / Consumer 用 bootstrap.servers 里随便找一台 Broker 连上去,发 MetadataRequest(topics)
  2. Broker 回包:返回该 Topic 所有 Partition 的元数据,包括每个 Partition 的 LeaderReplicasISR、以及对应 Broker 的 hostport
  3. 客户端建立 TCP 直连:客户端按 Partition Leader 的位置,和需要的 Broker 建立 TCP 连接(Producer 发到 Leader,Consumer 拉自 Leader)。
  4. 元数据失效再次刷新:触发条件包括:
    • metadata.max.age.ms 到期(默认 5 分钟)周期性刷新;
    • 收到 NOT_LEADER_FOR_PARTITION / LEADER_NOT_AVAILABLE 错误立刻重新 MetadataRequest
    • 收到新元数据版本(Broker 在响应里带 controller_epoch / metadata_version);
  5. 元数据传播: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 命令行家族 + Python confluent-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)

cluster_probe.py ↗