7146 字
36 分钟
Kafka 的本质:从消息队列到分布式事件日志

Kafka 的本质:从消息队列到分布式事件日志#

从同步调用、Queue 到 Partitioned Log,建立理解 Kafka 的第一性原理

阅读目标#

学完本篇,你不应该只会复述“Kafka 是一个高吞吐消息队列”,而应该能够从以下四个核心抽象重新推导 Kafka:

  • Log:事件不是等待被取走的任务,而是已经发生并被记录的事实;
  • Partition:单条 Log 无法无限扩展,因此需要把日志切分成多个独立有序分片;
  • Offset:消费者不是拥有消息,而是通过位置读取日志;
  • Consumer Position:事件数据与消费进度分离,同一份数据可以形成多个互不干扰的消费视角。

本文以 Apache Kafka 4.3 文档为基线。Kafka 官方将其定义为一个分布式事件流平台:它能够跨多台机器发布、读取、持久保存和处理事件。Kafka 4.x 的集群管理主线已经是 KRaft,而不是 ZooKeeper。1


0. 一条订单为什么拖垮了整个系统#

先不谈 Kafka。

假设我们正在开发一个订单系统。用户提交订单后,订单服务需要完成订单落库、库存锁定、风控检查、积分计算和通知发送:

public void createOrder(Order order) {
orderRepository.save(order);
inventoryService.lock(order);
riskService.check(order);
pointService.calculate(order);
notificationService.send(order);
}

在业务规模很小时,这段代码直观、简单,调用路径也容易追踪。但随着系统增长,问题会迅速出现:

  • 通知服务响应 3 秒,订单接口是否也必须多等待 3 秒?
  • 风控服务短暂不可用,用户是否就不能创建订单?
  • 积分服务处理能力只有订单服务的五分之一,峰值流量如何消化?
  • 如果以后增加推荐、数据仓库、审计和 Agent 平台,订单服务是否每次都要修改?

同步调用链中的耦合与故障传播

表面上看,这是一个“同步调用比较慢”的问题;但更深层的问题是:订单服务被迫知道订单创建之后整个世界需要做什么。

这意味着系统已经形成了多种耦合。


1. 同步调用真正的问题不是慢,而是耦合#

1.1 时间耦合:调用双方必须同时在线#

同步 RPC 的基本模式是:

A 发起请求 → B 处理请求 → B 返回响应 → A 继续执行

A 能够继续执行的前提是:

  1. B 此刻必须在线;
  2. B 此刻必须能够及时响应;
  3. A 必须在调用期间持续等待 B。

因此,A 与 B 被绑定在同一个时间窗口中,这就是 Temporal Coupling,时间耦合

以订单通知为例,“订单已经创建”和“短信已经发送”并不是同一个业务原子动作。订单已经成功落库后,即使短信平台暂时故障,订单事实本身也不应该消失。可是,在同步调用链中,订单请求却可能因为一个非核心下游服务而超时。

从业务语义看,我们真正希望表达的是:

订单已经创建。
至于谁关注这个事实、何时处理这个事实,由各自系统决定。

1.2 可用性耦合:下游故障会向上游传播#

假设调用链为:

Order → Inventory → Risk → Notification

只要其中一个强依赖不可用,整条链路就可能失败。即便每个服务都拥有较高可用性,串联之后的端到端可用性仍会下降。

更危险的是故障传播:

通知服务变慢
订单线程长时间等待
订单服务线程池被占满
新请求无法进入
故障从边缘服务扩散到核心交易链路

因此,同步依赖不仅连接了业务功能,也连接了不同服务的故障域。

1.3 扩缩容耦合:上游速度直接传递给下游#

假设订单服务突然收到 10,000 TPS:

10,000 次库存调用
10,000 次积分调用
10,000 次通知调用

如果通知系统每秒只能处理 2,000 个请求,它没有缓冲空间,只能通过超时、拒绝或排队把压力反向传给订单服务。

同步模型隐含了一个危险条件:

上游发送速度必须长期小于等于下游实时处理能力。

否则,速度较慢的一方会成为整条链路的瓶颈。

1.4 变化耦合:上游需要不断认识新的下游#

最初,订单服务只需要调用库存服务。后来业务不断增加:

Order Service
├── Inventory
├── Risk
├── Coupon
├── Points
├── Notification
├── Data Warehouse
├── Recommendation
└── Agent Platform

每新增一个“订单创建后的关注者”,订单服务都要增加依赖、修改代码、补充异常处理和重新发布。

此时,订单服务已经不再只负责订单,它还承担了业务流程编排中心的职责。

所以,大型系统引入消息机制的根本原因,并不只是“异步更快”,而是要解开:

事件发生者与事件关注者之间的时间、可用性、速度和变化关系。


2. 消息系统到底改变了什么#

2.1 从直接命令转向事实事件#

先区分两个容易混淆的概念:Command 与 Event。

Command:要求某个对象做一件事#

SendOrderCreatedSMS(orderId)
ReserveInventory(orderId)
CalculatePoints(orderId)

Command 通常包含明确的执行意图:谁来做、做什么。

Event:描述一件已经发生的事实#

OrderCreated(orderId)
PaymentCompleted(orderId)
OrderShipped(orderId)

Event 不要求某个具体系统执行动作,它只是公开一个不可否认的业务事实:

Order Service 发布 OrderCreated
┌────────┼─────────┬─────────┐
↓ ↓ ↓ ↓
Inventory Points Notification Analytics

订单服务只需要对“订单已经创建”这一事实负责。库存、积分、通知和数据分析系统分别解释这一事实对自己意味着什么。

这并不意味着所有服务都应该改成事件驱动,也不意味着使用 Kafka 就等于采用 Event Sourcing。这里真正重要的是:事件发生者不再需要知道全部事件消费者。

2.2 消息系统创建了一个时间缓冲区#

加入消息系统后,生产者与消费者不再要求同时工作:

Producer → Message Buffer → Consumer

假设生产速度短时达到 10,000 msg/s,而消费者只能处理 3,000 msg/s。消息系统可以先保存尚未处理的数据,消费者随后以自己的速度追赶。

因此,“削峰填谷”的本质并不是流量凭空消失,而是:

消息系统在两个速度不同的系统之间创建了一个时间缓冲区。

短期峰值被转化为积压;消费者获得时间处理积压;核心请求不必等待非核心动作同步完成。

不过,这里也必须保持警惕:缓冲区不能解决长期的生产消费能力倒挂。如果生产速度长期高于消费速度,积压仍会无限增长。Kafka 的消费者容量规划与 Lag 治理将在第七篇展开。

2.3 Broker 第一次出现#

最简单的消息架构是:

Producer → Broker → Consumer

Broker 位于生产者和消费者之间,负责接收、保存并提供消息。

但是,“Broker 中究竟保存什么”,决定了消息系统最终形成怎样的世界观。接下来,我们先看最符合直觉的 Queue。


3. Queue 是一个怎样的抽象#

3.1 Queue 把消息看成等待处理的任务#

最简单的队列模型是:

Producer: enqueue(message)
Consumer: dequeue()

消息从尾部进入,从头部被消费者取出。它所表达的核心语义是:

一项任务正在等待某个消费者处理。

消费流程通常被理解为:

Message → Deliver → Process → ACK → Completed

因此,Queue 主要关注的是 Message Delivery:

  • 消息有没有投递给消费者?
  • 消费者有没有确认处理?
  • 处理失败后是否重新投递?
  • 如何避免同一条消息同时交给多个竞争消费者?

3.2 多个消费者之后:竞争消费与发布订阅#

当多个消费者共同处理同一种任务时,可以形成竞争消费:

[M1][M2][M3][M4]
↓ ↓
Consumer A Consumer B

每条消息由其中一个消费者处理,用于提高任务处理并行度。

但如果库存、通知、风控和数据仓库都要接收 OrderCreated,就需要发布订阅:同一个事件需要被多个逻辑订阅者独立观察。

传统消息中间件同样可以支持这些模式,Kafka 也不是唯一实现方式。真正让 Kafka 的模型发生变化的,是它对另一个问题的回答。

3.3 为什么消息一定要“被取走”#

假设系统已经处理过以下事件:

10:00 OrderCreated
10:01 PaymentCompleted
10:05 OrderShipped

此时出现几个需求:

  • 数据仓库今天才上线,希望重新导入昨天的订单;
  • 风控算法升级,希望重新计算过去七天的风险;
  • 消费者代码存在 Bug,修复后需要重放过去一小时的数据;
  • 新增一个 Agent 分析服务,希望从历史事件中重建 Case 状态。

如果消息一旦成功消费就被理解为“任务完成并从系统中消失”,这些需求都需要额外的历史存储、复制链路或补数机制。

于是,一个更根本的问题出现了:

消息为什么必须被理解成等待领取的任务?它能不能被理解成已经发生、应当被记录的事实?

Queue 与 Log 的核心差异

Queue 的核心问题是:

这条消息应该交给谁?是否已经完成?

Log 的核心问题则是:

这份历史仍然存在;每个消费者读到哪里?

这就是理解 Kafka 的转折点。


4. Kafka 的核心世界观:Log#

4.1 Append-only Log#

最简单的 Log 是一个只能在尾部追加的有序记录序列:

Offset
0 1 2 3 4 5
│ │ │ │ │ │
[E0] → [E1] → [E2] → [E3] → [E4] → [E5] → append here

每次出现新事件,系统不是寻找一个可覆盖位置,也不是把旧记录弹出,而是在日志尾部追加新的 Record。

这种结构包含三个重要性质:

  1. 有序:记录按照写入顺序排列;
  2. 追加:新数据主要在尾部增加;
  3. 持久:消费者读取数据并不要求数据立即删除。

Kafka 官方设计中,一个 Topic 会被拆成多个有序 Partition;Consumer 在每个 Partition 上的位置可以用一个整数 Offset 表示,即下一条准备消费的记录位置。2

4.2 Queue 与 Log 的根本差异#

可以用一句话概括两者关注点:

Queue 的问题是“这条消息交给谁”,Log 的问题是“你读到哪里”。

Queue 往往围绕消息本身维护状态:等待投递、已投递、已确认、重新投递。

Log 则将数据状态与消费状态分离:日志只负责保存事件;消费者只负责记录自己在日志中的位置。

4.3 数据与消费进度分离#

假设日志包含 0~9 共十条事件:

0 1 2 3 4 5 6 7 8 9
───────────────────→

三个消费者可以拥有三个完全不同的位置:

Risk Consumer position = 8
Warehouse Consumer position = 5
Agent Consumer position = 2

同一份日志上的多个消费者位置

日志数据没有因为某个消费者读取而改变;一个消费者处理较慢,也不会迫使另一个消费者回退。

由此得到 Kafka 极其重要的抽象:

Event Data ≠ Consumer Progress

这看似只是把状态拆开,实际上改变了整个消息系统的能力边界。

4.4 Consumer 不拥有消息#

在典型 Queue 直觉中,消费者“拿走”了一条消息。

在 Log 直觉中,消费者只是发起如下读取:

fetch(partition=P0, offset=100)

Broker 返回从该位置开始的一批记录。消费者处理后,把自己的位置向前推进。

Kafka 官方设计强调,Consumer 可以控制自己的 Position,并能回退到旧 Offset 重新消费历史数据。2

因此,Consumer 与 Event 的关系不是所有权转移,而是:

消费者在一个持久有序的数据集合上移动读取位置。

4.5 Replay 不是附加功能,而是核心抽象的自然结果#

假设消费者已经读到 Offset 1000,后来发现 900~1000 的处理代码有 Bug:

0 ........ 900 ........ 1000
↑ ↑
replay from current

修复代码后,只要把消费位置重新调整到 900,就能再次读取这段事件。

Replay 来自两个基础条件:

Persistent Log
+
Movable Consumer Position

它并不是在传统 Queue 上额外添加的复杂补丁,而是 Log 模型自然产生的能力。

这体现了一个重要的架构原则:

优秀系统的高级能力,往往不是不断追加 Feature,而是由正确的核心抽象自然推导出来。

4.6 Kafka 应该如何定义#

到这里,我们可以先给 Kafka 一个底层定义:

Kafka 可以首先被理解为一个分布式、持久化、允许消费者按位置读取的事件日志系统。

这个定义用于建立 Mental Model,但它还不是 Kafka 产品能力的全部。Kafka 官方对它更完整的定义是“分布式事件流平台”,因为它不仅保存事件,还提供 Producer、Consumer、Streams、Connect、Admin 等能力。1

理解顺序应该是:先用 Log 看清底层抽象,再把 Kafka 看成构建在该抽象之上的完整 Event Streaming Platform。


5. 一条 Log 不够:为什么 Kafka 需要 Partition#

5.1 假设 Topic 只有一条全局 Log#

假设所有订单事件都写入一条日志:

Order Events
0 → 1 → 2 → 3 → 4 → ... → 100000000

所有 Producer 竞争同一个写入路径,所有数据存放在同一台机器,所有 Consumer 又希望保持全局顺序。

这个模型很快遇到四个限制。

写入吞吐限制#

所有写操作进入同一条 Log,写入并行度受限于单个日志分片所在机器的处理能力。

存储容量限制#

一条日志必须完整落在某个存储节点上,Topic 的容量上限受到单机磁盘约束。

网络带宽限制#

生产写入和消费读取都经过同一台机器,网络带宽难以横向扩展。

消费并行限制#

如果必须严格维持一条全局顺序,多个消费者就很难独立并行处理不同区段。

因此,仅有 Log 还不够。Kafka 还需要一个能够横向扩展的分片单位。

5.2 Partition 是独立有序的 Log#

Kafka 把一个 Topic 拆成多个 Partition:

Topic: order-events
├── Partition 0: 0 → 1 → 2 → 3
├── Partition 1: 0 → 1 → 2
└── Partition 2: 0 → 1 → 2 → 3 → 4

每个 Partition 都是一条独立有序日志,并拥有自己的 Offset 空间。

Topic、Partition 与 Offset

这里必须牢牢记住:

Offset 的作用域是 Partition,而不是整个 Topic。

正确坐标不是:

order-events 的 offset=100

而是:

TopicPartition(order-events, 0), offset=100

5.3 Partition 是 Kafka 的 Scale Unit#

Partition 可以分布到不同 Broker:

P0 → Broker 1
P1 → Broker 2
P2 → Broker 3

于是,一个 Topic 的数据和负载可以横向分布到多台服务器。

Kafka 官方运维文档明确指出,Partition 数决定 Topic 被切分为多少条 Log,也影响数据能够分布到多少服务器以及 Consumer 的最大并行度。3

因此,Partition 同时承担多个角色:

  • 存储分片单位:不同 Partition 能分布到不同 Broker;
  • 写入并行单位:Producer 可以同时写多个 Partition;
  • 读取并行单位:不同 Consumer 可以并行读取不同 Partition;
  • 顺序边界:Kafka 保证的是 Partition 内顺序,而不是 Topic 全局顺序。

5.4 Partition 不是免费的性能优化#

假设两个 Partition 分别包含:

P0: A → C → E
P1: B → D → F

我们只能确认:

P0 内 A 在 C 前,C 在 E 前;
P1 内 B 在 D 前,D 在 F 前。

但无法仅凭 Kafka Partition 顺序推导全局序列一定是:

A → B → C → D → E → F

因此,Partition 做出了一个架构交换:

放弃 Topic 全局顺序
获得水平扩展与并行能力

Partition 不是免费的性能优化,而是 Kafka 在“全局顺序”和“水平扩展”之间做出的架构选择。

5.5 Key 用于建立局部顺序#

实际业务往往不要求所有订单全局有序,只要求同一个订单的事件有序:

Order 1001 Created
Order 1001 Paid
Order 1001 Shipped

Producer 可以把 orderId 作为 Key:

hash(orderId) → Partition

相同 Key 通常会被映射到同一个 Partition,从而形成 Key 范围内的局部顺序:

Order1001 → P3
Order1001 → P3
Order1001 → P3

不过,Partition 数量变化可能改变 Key 到 Partition 的映射。Kafka 4.3 官方运维文档也提醒:如果分区逻辑依赖 hash(key) % partitionCount,增加 Partition 后,相同 Key 的后续消息可能进入不同 Partition,进而影响既有顺序假设。3

这一问题涉及 Topic 设计、Partition 扩容和容量治理,将在第七篇深入讨论。

5.6 Topic 到底是什么#

现在可以给 Topic 一个更准确的定义:

Topic 是一个逻辑事件流名称,其底层由一个或多个独立的 Partition Log 组成。

Topic 负责业务语义分类;Partition 负责物理切分、并行和局部顺序;Offset 负责定位 Partition 中的位置。

三者不能混为一谈。


6. Offset:Kafka 为什么不需要给每条消息维护消费状态#

6.1 Offset 首先是 Log Position#

在某个 Partition 中:

Offset 0 1 2 3 4 5
Record E0 E1 E2 E3 E4 E5

Offset 表示 Record 在该 Partition Log 中的位置。

它首先具有位置语义,而不是业务身份语义。

6.2 Offset 不等于 Message ID#

Message ID 主要用于唯一标识某条消息,例如 UUID:

550e8400-e29b-41d4-a716-446655440000

Offset 则表达一条记录在有序日志中的位置:

offset=1034

因为它是位置,所以天然支持:

从 1034 开始读取
读取 1000~1100 范围
回退到 900
计算距离日志尾部还有多少位置

因此:

Offset 首先是 Position,不应该被简单理解成业务消息 ID。

6.3 Consumer Position#

Consumer 读取 Partition 时,可以表达:

fetch(P0, offset=6)

Broker 从 P0 的 Offset 6 开始返回一批数据。Consumer 处理完毕后,下一次读取位置向前推进。

逻辑执行链为:

Current Position
Fetch Records
Process Records
Advance Position

Kafka 官方设计指出,Consumer 每个 Partition 的位置只需要一个整数,因此已消费状态非常小,并且可以周期性 Checkpoint。2

6.4 一条 Log 可以拥有多个消费视角#

对于同一个 order-events Partition:

Risk Consumer position=1000
Warehouse Consumer position=850
Agent Consumer position=200

它们可以:

  • 在不同时间上线;
  • 以不同速度处理;
  • 使用不同业务逻辑;
  • 在需要时分别回放;
  • 独立故障和恢复。

这就是数据与消费状态分离带来的独立性。

6.5 Current Position 与 Committed Offset#

本篇只做第一次区分。

Current Position#

当前运行中的 Consumer 已经读取到的位置,通常表示下一次 Fetch 将从哪里继续。

Committed Offset#

Consumer Group 持久记录的恢复位置。Consumer 崩溃或发生成员迁移后,新实例可以从已提交的位置恢复。

例如:

0 1 2 3 4 5 6 7 8 9
↑ ↑
committed current
offset=4 position=8

如果 Consumer 在处理到 8 后崩溃,但只提交到 4,那么恢复时可能重新读取 4 之后的数据。这正是交付语义、重复消费和 Offset 提交策略产生的地方。

不过,下面这些概念不能在第一篇混在一起:

Current Position
Committed Offset
Log End Offset
High Watermark
Last Stable Offset

第四篇将建立完整的 Offset 坐标系,第五篇再把它与副本、事务和可见性连接起来。

6.6 Kafka 重新定义了“消费”#

传统 Queue 经常把消费建模为消息状态变化:

NEW → DELIVERED → ACKED

Kafka 的经典 Consumer 模型则把消费建模为位置变化:

position=100 → position=120 → position=150

因此,可以得到全文最重要的结论之一:

Kafka 没有把“消费”建模成消息本身的状态变化,而是把它建模成消费者在 Log 上的位置变化。


7. 现在重新看 Kafka 的完整运行模型#

到这里,我们已经建立了五个概念:

Event → Log → Partition → Offset → Consumer Position

现在再把 Producer、Broker 和 Consumer 放回系统中。

Kafka 第一性原理运行模型

7.1 Producer:把 Event 追加到 Partition#

Producer 负责构造 Record,并将它发送到某个 Topic 的某个 Partition:

ProducerRecord
Topic
Partition Selection
Append to Partition Log

此处先不讨论序列化、Metadata、RecordAccumulator、Batch、Sender 和 NetworkClient。这些属于第三篇“Producer 完整执行链”。

第一篇只需要记住:

Producer 的最终目标不是“把消息交给某个 Consumer”,而是把 Event 追加到某条 Partition Log。

7.2 Broker:承担数据服务职责#

Broker 是承担 Kafka 数据服务职责的服务器角色,主要负责:

  • 接收 Producer 写入;
  • 保存 Topic Partition;
  • 响应 Consumer Fetch;
  • 参与 Partition 副本复制。

单个 Broker 可以保存多个 Topic 的多个 Partition;一个 Topic 的不同 Partition 也可以分布在不同 Broker 上。

7.3 Consumer:按 TopicPartition 和 Offset 读取#

Consumer 的核心行为可以抽象为:

Fetch(TopicPartition, Offset)

它从指定位置读取一批 Record,完成反序列化和业务处理,然后推进消费位置。

这解释了 Kafka 为什么更接近“读取日志”,而不是等待 Broker 主动把一条任务推到消费者手中。

7.4 用一句链路串起来#

Kafka 的基础数据链路可以表述为:

Producer
↓ append
Topic / Partition Log
↓ fetch(offset)
Consumer

其中:

  • Topic 提供业务语义;
  • Partition 提供分片、并行和局部顺序;
  • Offset 提供位置坐标;
  • Consumer Position 表示独立消费进度。

这四个抽象是后续理解 Kafka 高性能、Producer、Consumer Group、副本和事务的共同地基。


8. Kafka 是一个分布式系统:Broker、Controller 与 KRaft#

8.1 为什么单 Broker 不够#

单 Broker 面临明显边界:

  • 单机磁盘容量有限;
  • 单机网络和 CPU 吞吐有限;
  • Broker 故障可能导致服务不可用;
  • 数据只保存一份无法容忍磁盘或节点故障。

因此,Kafka 以 Cluster 形式运行,让 Partition 分布到多个 Broker,并为 Partition 建立副本。

此处只建立概念:

Partition
├── Leader Replica
└── Follower Replica(s)

Producer 和 Consumer 的正常数据读写围绕 Partition Leader 进行;Follower 复制 Leader 数据。Replica、ISR、High Watermark 和 Leader 故障恢复将在第五篇深入研究。

8.2 数据之外还有集群元数据#

Kafka Cluster 必须持续知道:

  • 集群中有哪些 Broker;
  • 存在哪些 Topic;
  • Topic 包含哪些 Partition;
  • Partition Replica 位于哪些 Broker;
  • 当前 Leader 是谁;
  • 配置和权限如何变化。

这些信息属于 Cluster Metadata。

如果没有一致的 Metadata,Producer 无法知道应向哪个 Broker 写入,Consumer 也无法找到 Partition Leader。

8.3 KRaft:现代 Kafka 的控制平面#

Kafka 4.x 的架构主线是 KRaft。服务器可以通过 process.roles 配置为:

broker
controller
broker,controller

Kafka 官方文档说明,Controller 节点参与 Metadata Quorum,其中一个是 Active Controller,其余作为 Hot Standby;生产环境通常将 Controller 与 Broker 职责隔离。4

KRaft 控制平面与 Broker 数据平面

可以借用 Control Plane / Data Plane 理解:

Data Plane:Broker#

处理业务数据路径:

Producer Write
Consumer Fetch
Partition Storage
Replica Transfer

Control Plane:KRaft Controller Quorum#

管理集群控制状态:

Topic / Partition Metadata
Leader Information
Broker Registration
Configuration Changes

这种划分并不意味着 Controller 不重要,恰恰相反:Controller Quorum 决定集群是否能持续完成元数据变更和故障管理。

8.4 为什么本专栏不再以 ZooKeeper 为主线#

Kafka 早期版本依赖 ZooKeeper 保存和协调部分集群元数据,许多旧教程仍然从 ZooKeeper 开始讲 Kafka。

但学习现代 Kafka 时,应以 KRaft 为默认架构。ZooKeeper 只在解释历史演进、旧集群迁移或阅读旧源码时作为背景出现。

这可以避免形成一种已经过时的认知:

Kafka = Broker + ZooKeeper

现代主线应该是:

Kafka Cluster
├── Broker Data Plane
└── KRaft Controller Quorum

9. Kafka 到底是不是消息队列#

回答这个问题时,不应该走向两个极端。

错误说法一:

Kafka 不是 MQ。

Kafka 显然可以承担大量异步消息、发布订阅和任务解耦场景。

错误说法二:

Kafka 就是一个普通 MQ。

只用 Queue 模型无法完整解释 Kafka 的历史保留、Replay、独立 Consumer Position、流处理和数据管道能力。

更准确的说法是:

Kafka 可以承担消息队列场景,但只有使用 Partitioned Log Mental Model,才能完整理解它。

Queue Mental Model:

Producer → Queue → Consumer

Kafka Mental Model:

Producer
↓ Append Event
Partitioned Persistent Log
↓ Read by Position
Consumer

Kafka 官方将其定位为分布式事件流平台,因为它组合了事件发布订阅、持久保存和事件处理能力。1


10. 用四个抽象重新推导 Kafka#

本文没有从 API 开始,也没有从安装命令开始,而是从系统矛盾一步步推导出 Kafka。

10.1 从 Log 推导#

Log
├── 事件有序追加
├── 历史可以保留
├── 消费不要求删除数据
└── Consumer 可以 Replay

10.2 从 Partition 推导#

Partition
├── 分布式存储
├── 并行写入
├── 并行读取
├── Partition 内顺序
└── 放弃 Topic 全局顺序

10.3 从 Offset 推导#

Offset
├── 定位 Partition 中的位置
├── 支持范围读取
├── 支持回退和重放
└── 支持恢复消费进度

10.4 从独立 Consumer Position 推导#

Independent Consumer Position
├── 多个系统读取同一份数据
├── 不同消费者拥有不同速度
├── 消费者故障不改变日志
└── 每个消费者可以独立 Replay

最终,Kafka 的底层模型可以收束为:

Persistent Log
+
Partitioning
+
Offset-based Reading
+
Independent Consumer Progress
Distributed Event Streaming Platform

Kafka 最重要的设计,并不是创建了 Producer 和 Consumer API。

它真正改变的是我们对“消息”的理解。

在传统 Queue 直觉中,一条消息是一项等待被交付和完成的任务;在 Kafka 的 Log 世界中,一条 Event 是已经发生并被记录的事实。

Consumer 不再等待 Kafka 把消息交给自己。

它只是站在一条持续增长的 Log 上,决定自己要从哪里开始阅读。


11. 第一篇自测与验收#

完成本篇后,应当能够脱稿回答以下问题:

  1. 为什么同步 RPC 的核心问题不只是延迟?
  2. Temporal Coupling、Availability Coupling、Scaling Coupling 和 Change Coupling 分别是什么?
  3. 为什么消息系统可以被理解为时间缓冲区?
  4. Command 和 Event 的语义有什么不同?
  5. Queue 主要围绕什么问题建立状态?
  6. Queue 与 Log Mental Model 的根本差异是什么?
  7. Kafka 的 Replay 为什么是核心抽象自然产生的能力?
  8. 为什么只有一条 Log 无法支撑 Kafka 水平扩展?
  9. Partition 同时带来了哪些能力和代价?
  10. 为什么 Offset 的作用域必须是 TopicPartition?
  11. Offset 为什么不能简单等价于 Message ID?
  12. Current Position 与 Committed Offset 有什么区别?
  13. Broker Data Plane 与 KRaft Control Plane 分别负责什么?
  14. 为什么只用“消息队列”无法完整理解 Kafka?

如果只能背出 Topic、Partition、Producer、Consumer 的定义,本篇还没有真正学会。

如果可以从 Log、Partition、Offset 和 Consumer Position 推导出 Kafka 的历史保留、Replay、扩展性、局部顺序和独立消费,你已经建立了正确的 Kafka 世界观。


12. 下一篇预告:Kafka 为什么快#

现在我们已经把 Kafka 理解成一组持续增长的 Partition Log。

接下来会出现更底层的问题:

  • Partition 在磁盘上到底是什么?
  • 为什么 Partition 还要拆成多个 Segment?
  • .log.index.timeindex 分别解决什么问题?
  • Consumer 请求 Offset 36891 时,Broker 如何定位磁盘位置?
  • 顺序 I/O、Page Cache、Batch、Compression 和数据传输路径如何共同构成 Kafka 的高吞吐?
  • Kafka 的性能究竟来自单个技巧,还是来自整个 Data Path 的协同设计?

第二篇将从 Kafka 的磁盘文件开始,完整拆解:

《Kafka 为什么快:从 Partition Log、Segment 到 Page Cache 的数据路径》


参考资料#

  1. Apache Kafka Documentation / Introduction
  2. Apache Kafka 4.3 Design
  3. Apache Kafka 4.3 Basic Kafka Operations
  4. Apache Kafka 4.3 KRaft
Kafka 的本质:从消息队列到分布式事件日志
https://jupiter-ws.cn/posts/backend/kafka/01_kafka_essence_from_queue_to_distributed_event_log/
作者
Jupiter
发布于
2026-07-16
许可协议
CC BY-NC-SA 4.0