Chapter 01

核心概念:把 Kafka 当成一条日志

起点页给出了一句话本质——Kafka 不是消息队列,而是一条分布式、可重放、按分区切分的提交日志(commit log)。这一章把这句话拆成六个互相咬合的概念,让"日志"这个心智模型真正落地,不再是一句口号。

本章你将建立的 schema

  • 一条 partition 是 append-only 的日志,每条记录有单调递增的 offset;消费 = 移动游标,读完不删除。
  • 顺序只在单个 partition 内成立;topic 是逻辑分类,并行度上限 = partition 数。
  • offset 归消费者维护,不归 broker;不同 consumer group 各读各的同一份日志。
  • producer 按 key 哈希选分区;broker 持有 partition 的 leader/follower 副本。

1.1为什么需要 Kafka:队列删一次,日志读多次

传统 MQ 把消息当"消费一次就丢"的任务;Kafka 把消息当"写进日志、谁都能重读"的事实。

为什么需要它

传统消息队列(RabbitMQ、ActiveMQ 这类)的核心动作是投递并删除:一条消息被某个消费者 ack 之后,broker 就把它从队列里移除,消费进度由 broker 记账。这套模型服务于"任务分发"——一封邮件发一次、一笔扣款扣一次。

但当同一份数据要被多个下游各自独立消费时,这个模型就崩了:风控要读订单流、数仓要读订单流、推荐要读订单流,三方进度不同、还会回溯重算。在删除式队列里,要么给每个下游复制一份队列,要么消息删早了导致后来者读不到。没有 Kafka,工程师得自己搭一套"留存原始事件 + 各下游记自己读到哪"的基础设施——而这恰好就是一条带游标的日志。

底层机制(比文档深一层)

这两类系统的分水岭不在"快慢",而在谁持有消费进度、消息何时消失。传统 MQ 把"已读到哪"作为 broker 的内部状态,消息的生命周期绑定在"是否被 ack"上——ack 即删除,进度无法回退。Kafka 反过来:broker 只负责把记录顺序追加进日志、按时间或大小(而非"是否被消费")做保留,消费进度 offset 是消费者自己提交的一个数字。把"进度"从 broker 搬到消费者这一侧,是后面一切特性的总开关:消息不再因被读而消失,于是可重放;多个消费者各存各的 offset,于是同一份数据能被独立消费多次。

类比 · 带边界声明

传统队列像取号机的叫号小票:叫到你,小票作废,下一个人看不到你那张。Kafka 像报纸的合订本:今天的报纸印出来摆上架,张三李四都能翻、还能翻回上周。边界:报纸合订本会无限堆下去,Kafka 不会——它按保留期(默认 7 天)或容量删旧日志段,过期的报纸会被回收。所以 Kafka 是"有保留窗口的可重放",不是"永久存储"。

场景走查

一个订单系统每秒产生几千条"订单已创建"事件。用删除式队列:风控消费完一条就被删,数仓再想读同一条已经没了,只能各开一个队列、生产者发三遍。用 Kafka:事件写进 orders 这一条日志一次,风控、数仓、推荐三个独立消费组各自维护 offset、各读各的;某天推荐算法改版要重算上周数据,把它的 offset 重置到 7 天前再跑一遍即可——日志还在,重放不需要生产者配合。

与下一个概念的关系:把"进度归消费者、读完不删"这件事讲到底,就必须看清这条"日志"在物理上长什么样——它就是下一节的提交日志模型。

1.2提交日志模型:append-only + offset

记录只能追加到日志末尾、永不修改,每条带一个单调递增的 offset;消费就是按 offset 顺序移动一个读游标。

为什么需要它

"提交日志"不是 Kafka 发明的——数据库的预写日志(WAL)、Raft 的复制日志都是同一个结构:一串只追加、不可改、严格有序的记录。它之所以是分布式系统的地基,是因为"一串有序且不可变的记录"是世界上最容易被复制、被重放、被多方达成一致的数据结构。两台机器只要从头到尾按相同顺序回放同一条日志,状态就必然一致。Kafka 把这个结构直接暴露成产品。

底层机制(比文档深一层)

offset 不是"消息 ID",而是记录在这条日志里的位置序号——从 0 开始,每追加一条 +1,分区内永不重复、永不回退。这个设计带来三个直接后果:

  • 写入是 O(1) 顺序追加:永远只往 active 日志段的尾部写,不需要像 B 树那样随机寻址、加锁、再平衡。顺序磁盘 I/O 接近内存速度,这是 Kafka 高吞吐的物理根源(机制细节见 §2.1 存储)。
  • 读取是"从 offset N 开始往后给记录":消费者发来一个 offset,broker 从那个位置往后顺序吐数据。读不破坏写、也不互相干扰,因为各消费者只是停在日志不同位置的游标。
  • 读完不删:游标前移不影响日志本身。记录何时消失只取决于保留策略,与"是否被读过"完全解耦。

这正是起点页那句"门槛"心智模型的落点:消费不是出队(dequeue),是移动游标(seek)。一旦接受这一点,"为什么消息读完还在""为什么能从头重放""为什么 offset 在消费者手里"就不再是需要单独记忆的知识点,而是同一个结构的必然推论。

LOG v0 0 v1 1 v2 2 v3 3 v4 4 v5 5 v6 6 下一条 7 Producer 追加 组 A · offset 5 组 B · offset 2
图 1.1同一条日志,producer 在尾部追加,两个消费组的游标各停一处。 注意:组 A 已读到 5、组 B 才到 2,但 0–4 号记录仍在日志里——offset 是消费者各自的游标,不是"已删除到哪"。
类比 · 带边界声明

日志像账本:只在最后一页往下记,写错了不能擦,只能再记一笔冲正。offset 像页码加行号,"读到第几行"是每个读者自己拿书签夹着的。边界:账本不会被删,Kafka 日志会按保留期截断——书签所指的那一行一旦超过保留期就被撕掉(消费滞后超过保留期即丢数据,是 04 章的一类失败模式)。

场景走查

消费者 poll 一批记录、处理完、把 offset 5 提交到 broker;进程崩溃重启后,它向 broker 要"从 offset 5 之后的记录",于是从 6 继续,不重不漏。如果它在崩溃前没来得及提交 offset,重启会从上次提交的位置(比如 3)重新拉,4、5 被重复处理——这把"提交时机"变成正确性问题,留到 04 章展开。

想一想

两个不同的 consumer group 读同一个 topic。组 A 已经读到 offset 100,组 B 才读到 offset 10。组 B 会因为"落后"而漏掉中间的消息吗?

展开答案(先停 10 秒再点)

不会。offset 是每个组各自维护的游标,互不影响。组 A 读到 100 不会"消耗"掉记录——11 到 100 号记录全都还在日志里(只要没超过保留期),组 B 会照常从 11 一路读下去。

这道题指向的设计要点:删除式队列里"别人读了你就没了",而 Kafka 里读取只是移动自己的游标,多个消费组天然隔离。这正是"日志而非队列"带来的、最反直觉也最有用的一条性质。

与下一个概念的关系:到这里"日志"还是一条。但单条日志只有一个写入点、无法水平扩展。下一节看 Kafka 怎么把一条逻辑日志切成多条物理日志——这就是 topic / partition / offset。

1.3Topic / Partition / Offset:顺序只在分区内

topic 是逻辑分类,物理上切成 N 个 partition,每个 partition 是一条独立有序的日志,offset 是分区内的位置。

为什么需要它

如果一个 topic 只有一条日志,那它只有一个写入尾部、一个读取序列——吞吐被单机磁盘和单消费者卡死,无法水平扩展。partition 是 Kafka 的解法:把一个 topic 切成多条并列的日志,分散到不同 broker 上,写入和消费就能并行。partition 是顺序和并行的共同单位——这一句要记牢,后面的取舍全从它来。

底层机制(比文档深一层)

关键的、面试最爱考的一点:Kafka 只保证单个 partition 内有序,不保证 topic 全局有序。原因是物理的——全局有序需要一个单一写入点把所有记录串成一条序列,那就等于退回单分区、放弃了扩展性。所以 offset 是分区内的位置序号:partition 0 有它自己的 0,1,2…,partition 1 也有自己的 0,1,2…,两个分区的"offset 5"毫无关系。一条记录的完整坐标是 (topic, partition, offset) 三元组,而不是单个 offset。

这意味着:需要保证先后顺序的记录(比如同一个订单的"创建→支付→发货")必须落进同一个 partition,否则它们分散在不同分区、消费端无法保证读到的先后。怎么让它们落同一分区?靠 key——这是 §1.4 的事。

类比 · 带边界声明

topic 像一条高速公路,partition 像这条路上的多条车道:同一车道内车辆前后有序,但你不能断言"3 号车道的第 5 辆车"比"1 号车道的第 5 辆车"先出发。边界:真实车道之间可以变道,Kafka 的记录一旦按 key 进了某条 partition 就不会跨分区移动;而且 partition 数只能增不能减,加分区还会打乱已有的 key→分区映射(代价见 §2.2)。

场景走查

orders topic 切成 4 个 partition。生产者用 orderId 作 key,于是 order-42 的所有事件(创建、支付、发货)都哈希到同一个 partition,消费端读这个分区时它们必然按写入顺序到达。而 order-42 和 order-99 可能落在不同分区——它们之间没有顺序保证,这通常也无所谓,因为两笔订单本就互不相关。"只在需要顺序的范围内保证顺序"是 Kafka 的核心权衡。

与下一个概念的关系:分区解决了"日志怎么切",但还没说"谁来写、谁来读、怎么把 N 个分区分给多个消费者并行处理"。这是 producer / consumer / consumer group。

1.4Producer / Consumer / Consumer Group:并行单位 = 分区数

producer 按 key 选分区写入;一个 consumer group 内每个 partition 只分给一个 consumer;不同 group 各自独立消费同一份数据。

为什么需要它

有了多个 partition,就需要一种机制把它们分配给多个消费进程并行处理,同时还得保证"同一分区不被组内两个消费者同时读"——否则分区内的顺序保证就被两个消费者撕碎了。consumer group 就是这个分配机制:组内成员瓜分分区,组间互不干扰。它让"扩消费能力"变成"往组里加消费者"这么简单——但有个硬上限。

底层机制(比文档深一层)

分配的铁律:一个 partition 在同一时刻只能被同一个 consumer group 里的一个 consumer 消费。由此推出 Kafka 并行度的硬上限——一个组的有效并行度 = partition 数。组里消费者比分区多,多出来的就空闲拿不到分区;比分区少,则有消费者要扛多个分区。这是规划分区数时第一个要算的约束。

"谁拿哪个分区"由一个叫再平衡(rebalance)的过程决定:成员加入/退出时重新分配分区。再平衡的代价、协议演进(eager 急切式 vs cooperative 协作增量式)是 §2.4 的重头戏,这里只需知道它存在、它在成员变动时触发。另一条正交的线:不同的 consumer group 读同一个 topic 完全独立——各存各的 offset、各按各的进度,互不影响(图 1.1 已画过两个组停在不同 offset)。

TOPIC orders partition 0 日志 · offset 0,1,2… partition 1 日志 · offset 0,1,2… partition 2 日志 · offset 0,1,2… GROUP 风控 consumer C1 consumer C2 consumer C3 分给 3 分区 → 3 消费者,1:1 占满
图 1.2组内每个 partition 恰好连一个 consumer。 注意:连线是一对一——这正是"并行度 = 分区数"的来源。再加第 4 个消费者进这个组,它会拿不到分区而空闲。

一个 ProducerRecord 长什么样

生产者发出去的不是一个裸字符串,而是一个带 key 的 ProducerRecord。key 决定它落进哪个分区——这是把"需要顺序的记录"钉在同一分区的方式:

ProducerRecord 锚点片段 Java
// 第 1 个参数 = topic,第 2 个 = key,第 3 个 = value
// key = orderId:同一订单的所有事件哈希到同一 partition,于是有序
var record = new ProducerRecord<>("orders", order.getId(), order.toJson());
producer.send(record);   // 异步追加到该 key 对应 partition 的日志尾部
key不是数据库主键,而是路由依据:相同 key → 相同 partition → 同一条日志 → 有序。key 传 null 时记录在分区间均摊(见 §1.6),就失去按 key 的顺序保证。
记录(key) 分区器 PARTITION key = order-42 key = order-99 key = order-42 murmur2(key) % 分区数 partition 0 partition 1 partition 2 同 key 同分区
图 1.3key 经 murmur2 哈希取模选定分区。 注意:两条 order-42(朱红)无论何时发,都落进同一个 partition 1——这就是"相同 key → 同分区 → 有序"的来源;key 不同(order-99)则可能落到别的分区。

消费侧的 poll 循环长什么样

消费者不是被 broker"推"消息,而是自己循环 poll(拉)。这把消费节奏的控制权交给消费者,是它能管自己 offset 的前提:

consumer poll 锚点片段 Java
consumer.subscribe(List.of("orders"));   // 加入消费组,由再平衡分到若干 partition
while (running) {
    var records = consumer.poll(Duration.ofMillis(500));  // 主动拉一批
    for (var r : records) {
        handle(r.value());                // 处理记录
    }
    consumer.commitSync();                // 处理完再提交 offset(先处理后提交 = 至少一次)
}
poll既拉数据也"证明存活"——长时间不调 poll 会被判死并触发再平衡(04 章的再平衡风暴)。commitSync放在处理之后,决定了投递语义;放处理之前会把失败模式从"重复"翻成"丢失"(§2.5)。
想一想

一个 topic 有 4 个 partition。你给同一个 consumer group 启动了 6 个 consumer。会发生什么?

展开答案(先停 10 秒再点)

4 个 consumer 各分到 1 个 partition,剩下 2 个 consumer 完全空闲,一条记录都拿不到。因为"一个 partition 同一时刻只能给组内一个 consumer",组的并行度被 partition 数顶死在 4,加再多消费者也不会提升吞吐。

设计要点:扩消费能力的前提是分区数足够。这也解释了一个反直觉现象——盲目加消费者不仅无效,每次加入还触发一次再平衡,反而可能让 lag 变大(04 章)。规划分区数时,先想清楚目标消费并行度。

与下一个概念的关系:producer 和 consumer 都在跟"partition"打交道,但 partition 的副本到底存在哪台机器上、读写打到哪个副本?这是 broker / cluster / replica。

1.5Broker / Cluster / Replica:日志被复制到多台机器

broker 是 Kafka 服务进程,一个 partition 有 1 个 leader + 若干 follower 副本,分布在不同 broker 上。

为什么需要它

partition 是一条日志,落在某台 broker 的磁盘上。如果只存一份,这台机器一挂,这个分区的数据就没了、也读不了。replica(副本)就是把同一条 partition 日志在多台 broker 上各存一份,让单机故障不导致数据丢失或分区不可用。一组 broker 协同工作就是一个 cluster(集群)。

底层机制(比文档深一层)

同一个 partition 的多个副本里,有且只有一个是 leader,其余是 follower。机制的关键点:

  • 所有读写都只打 leader,follower 不直接对外服务——它们像消费者一样主动从 leader拉取(fetch)新记录,把 leader 的日志复制到自己这里。这保证了所有副本回放的是同一条日志、顺序一致。
  • 副本分布在不同 broker 上,所以一台 broker 宕机时,它上面那些 partition 的 leader 角色会切换到别的 broker 上的 follower,服务继续。
  • follower 并非总能跟上 leader。Kafka 用一个叫 ISR(同步副本集)的集合追踪"哪些副本目前跟得够紧",并用 HW(高水位)控制消费者能读到哪一条——这两个机制决定了"持久性"和"消费者可见性"的边界,是 §2.3 的核心,详见 02 章,本章不展开。
类比 · 带边界声明

leader/follower 像一个主记账员 + 几个抄录员:所有人只把账报给主记账员记,抄录员照着主账本一行行抄一份备份。主记账员请假,立刻从抄录员里指定一个接任。边界:现实里抄录员可能抄得比主账本还全,但 Kafka 的 follower永远不会领先 leader(它只能拉已写进 leader 的记录);而且能不能接任、接任会不会丢记录,取决于这个 follower 当时是否在 ISR 里——这正是 02 章要讲的 unclean leader election 风险。

场景走查

orders 的 partition 0 配了 3 个副本(replication.factor=3),分别在 broker-1、broker-2、broker-3,leader 在 broker-1。生产者发往 partition 0 的记录都写到 broker-1,broker-2/3 上的 follower 持续 fetch 跟上。broker-1 突然宕机:集群从 broker-2/3 的 follower 里选一个当新 leader,生产和消费切到新 leader 继续——读者侧几乎无感。这就是副本机制:它复制的不是别的,正是那条 partition 日志。

与下一个概念的关系:从宏观(集群、副本)回到微观——日志里那一条"记录"本身到底由什么组成?key 除了路由还管什么?这是最后一个概念 Record。

1.6Record:key / value / headers / timestamp

日志里的一条记录由 key、value、headers、timestamp 组成;key 同时决定分区路由与 compaction 行为。

为什么需要它

前面五节都在说"记录",这一节把它拆开。把 record 理解成"只是个 value"会错过 Kafka 最巧的一处设计:key 不是可有可无的标签,而是同时控制两个关键行为——它落哪个分区、以及在日志压缩里它代表"哪个实体的最新状态"。

底层机制(比文档深一层)

一条 record 的四个部分:

  • value:消息体本身(订单 JSON、事件 payload)。
  • key:路由依据。有 key → murmur2 哈希到固定分区,相同 key 永远同分区(顺序保证的来源);key 为 null → 由 sticky partitioner 在分区间均摊(机制见 §2.2)。
  • headers:键值对元数据(如 trace-id、schema 版本、来源系统),不影响路由,供消费端按需读取。
  • timestamp:记录的时间戳(生产时间或日志追加时间),是按时间检索、按时间保留的依据。

key 的第二重身份在日志压缩(log compaction)里:开启压缩的 topic 会对每个 key 只保留最新那条 value,旧值被回收;发一条 value 为 null 的记录(称为 tombstone 墓碑)则表示"这个 key 删除了"。于是一个压缩 topic 就成了"每个 key 的最新状态"的快照——这是 Kafka 能给一张表当 changelog、能存消费组 offset 的底层原理。压缩的具体机制(何时触发、tombstone 何时清)留给 §2.1。

类比 · 带边界声明

普通 topic 像流水账(每笔都留),压缩 topic 像余额表(每个账户只留最新余额)——而决定"同一个账户"的正是 key。边界:压缩是最终去重,不是即时——旧值会在后台清理触发前一直存在,tombstone 也有保留窗口。所以不能假设"写了新值旧值马上消失"(细节见 02 章)。

场景走查

用户资料 topic 以 userId 为 key、开启 compaction。用户改了三次昵称,日志里一度有三条记录;压缩后只剩最新那条。新启动的服务从头读这个 topic,就能重建"每个用户的当前资料"全量快照,不必读完整个历史。注销用户时发一条 userId + null 的 tombstone,压缩后这个 key 彻底消失。同一个 key 既决定路由、又定义"同一实体"——这就是 key 的双重身份。

§本章 self-check

先合上教程,把你能想到的答案写在纸上或编辑器里。 写完再点开答案对照——直接点开等于把这一节当再读一遍。

  1. 用一句话说清 Kafka 的"消费"和传统消息队列的"消费"在消息生命周期上的根本区别。
  2. offset 是全局唯一的吗?一条记录的完整坐标由哪几部分组成?为什么 Kafka 不保证 topic 全局有序?
  3. 一个 consumer group 的最大有效并行度由什么决定?为什么往组里无限加 consumer 不能无限提升吞吐?
  4. (设计题)你要设计一个 topic,承载"用户余额变更"事件,要求:① 同一用户的变更严格有序;② 新服务启动时能快速重建每个用户的当前余额,不必回放全部历史。你会怎么选 key?topic 用普通保留还是开 compaction?为什么?
答案(先做完再展开)
  1. 传统队列:消息被 ack 后从 broker 删除,消费进度由 broker 记账,读完即消失、不可重放。Kafka:消息追加进日志后按时间/大小保留、不因被读而删除,消费只是移动消费者自己维护的 offset 游标,因此可被多个消费组独立重读。
  2. 不是全局唯一。offset 只在单个 partition 内单调递增,完整坐标是 (topic, partition, offset) 三元组。不保证全局有序是因为全局有序需要单一写入点串行化所有记录,等于退回单分区、放弃水平扩展——Kafka 选择"只在分区内有序"换取扩展性。
  3. 由该 topic 的 partition 数决定。因为一个 partition 在同一时刻只能被组内一个 consumer 消费,consumer 数超过 partition 数时多出来的只能空闲,所以并行度顶死在分区数;加 consumer 还会触发再平衡,可能适得其反。
  4. key 选 userId——保证同一用户所有变更哈希到同一 partition,从而分区内有序(满足 ①)。topic 开 compaction(日志压缩)——每个 userId 只保留最新一条 value,新服务从头读即可重建"每个用户当前余额"的快照,无需回放全部历史(满足 ②)。注销用户用 userId + null 的 tombstone 删除。这道题的判别点:需要"按实体最新状态"时用 compaction(changelog 语义),需要完整事件流时用普通时间保留——两者由 topic 配置区分,且都依赖 key 定义"同一实体"。
进阶挑战 · 刚好够不着

如果业务要求"全局严格顺序",分区数该怎么定?

设想一个场景:一个审计系统要求整个 topic 的所有事件都严格按写入先后被消费,不允许任何两条记录乱序(不是按 key,是全局)。结合本章"顺序只在 partition 内"和"并行度 = 分区数"两条事实,推一推:这个 topic 的分区数只能是多少?它会牺牲掉 Kafka 的什么能力?如果业务量大到单分区扛不住,这个"全局严格顺序"的需求本身是不是哪里有问题?

提示(卡住再展开)

全局有序 ⇒ 只能有一条有序日志 ⇒ 分区数只能是 1。代价:消费并行度被锁死为 1(一个组只有一个 consumer 干活)、吞吐被单分区单机顶死,Kafka 的水平扩展全部失效。这通常是个信号——"真的需要全局顺序,还是只需要按某个 key 的顺序?"把全局顺序拆成按 key 的顺序,往往是更对的设计。这个取舍正是 §2.2 分区与顺序的核心,02 章会算这笔账。