阶段一 · 基础与核心模型

核心概念:Broker,Topic,Partition,Offset

一句话总结

Kafka 的数据模型是一个只追加的分布式日志:消息写进 Topic 的某个 Partition,每条消息有唯一递增的 Offset;多个 Broker 组成集群对外提供读写.
把握住"append-only log + offset 寻址"这个心智模型,后续所有机制都是它的衍生.

用一条消息的旅程串起全部概念

假设你的 Go 服务要把一条订单事件发到 Kafka,消费端读取后写入数据仓库.整个流程涉及的核心角色:

Producer(Go 服务)
│
│ 指定 Topic = "order-events"
│ 消息 Key = "user-123"
▼
┌────────────── Kafka Cluster ──────────────┐
│ │
│ Broker 0 Broker 1 Broker 2 │
│ ┌──────────┐ ┌──────────┐ │
│ │ Partition0│ │ Partition1│ ... │
│ │ offset 0 │ │ offset 0 │ │
│ │ offset 1 │ │ offset 1 │ │
│ │ offset 2 │ │ offset 2 │ │
│ │ ... │ │ ... │ │
│ └──────────┘ └──────────┘ │
│ │
└───────────────────────────────────────────┘
│
│ Consumer 从 Partition0, offset=2 开始拉取
▼
Consumer(ETL 服务)

Broker:集群中的一台服务器

Broker 是 Kafka 集群里的一个进程实例(通常一台机器跑一个).它的职责:

  • 接收 Producer 写入的消息,持久化到本地磁盘
  • 响应 Consumer 的拉取请求
  • 与其他 Broker 同步副本数据
属性 说明
标识 每个 Broker 有唯一的数字 ID(broker.id)
无状态(对客户端) 不记录"消费到哪了",offset 由 Consumer 自己管理
可水平扩展 加 Broker → 分区可以迁移/新建到新节点上

类比 Go 里的 goroutine worker pool:每个 Broker 是一个 worker,Topic 的 Partition 是任务,Kafka 把 Partition 分配给不同 worker 处理.

常见误解

Broker 不是"消息路由器"——它不做转发.Producer 直接把消息写到目标 Partition 所在的 Broker(通过 metadata 知道谁负责哪个 Partition).


Topic:逻辑上的消息分类

Topic 是一个逻辑命名空间,类似数据库里的表名.所有关于同一业务领域的消息归到同一个 Topic.

// 你的 Go 代码里就是一个字符串
topic := "order-events"

Topic 本身只是名字,不存储数据.数据存在 Partition 里.

命名建议:

模式 示例 说明
<领域>.<事件类型> order.created 事件驱动风格
<服务名>.<实体> payment-svc.transactions 服务导向风格
<环境>.<领域>.<事件> prod.order.created 多环境共用集群时

Partition:并行与顺序性的基本单位

这是 Kafka 最核心的抽象.一个 Topic 被切分为 N 个 Partition,每个 Partition 是一个有序的、不可变的消息序列(append-only log).

为什么要分区

单个日志文件的读写吞吐有上限.把一个 Topic 拆成多个 Partition,就能:

目的 原理
并行度 不同 Partition 可以分布在不同 Broker 上,读写并行
水平扩展 单个 Partition 有吞吐上限(通常 10-50 MB/s),多分区叠加
消费并行 一个 Consumer Group 里,每个 Consumer 分配若干 Partition,并行消费

Partition 与 Broker 的关系:N 对 1

一个 Broker 上可以承载很多个 Partition(来自不同 Topic,甚至同一 Topic 的多个分区).Kafka 的 Controller 会尽量把同一 Topic 的 Partition Leader 均匀分散到不同 Broker,实现负载均衡:

3 Brokers,Topic "order-events" 有 6 个 Partition(副本因子=2):

Broker 0: P0(leader), P3(leader) | P1(follower), P4(follower)
Broker 1: P1(leader), P4(leader) | P2(follower), P5(follower)
Broker 2: P2(leader), P5(leader) | P0(follower), P3(follower)

每个 Broker 当了 2 个 Partition 的 Leader(负责读写),
同时存了另外 2 个 Partition 的副本(只做同步备份).

关键结论:Partition 是负载分散的单位.6 个 Partition 分布在 3 个 Broker → 读写压力被 3 台机器分摊.

消息路由:Key 决定去哪个 Partition

Producer 发消息时可以附带一个 Key(任意 []byte),Kafka 根据 Key 决定消息写入哪个 Partition:

// 设置 Key → hash 路由,同 Key 的消息落同一 Partition
record := &kgo.Record{
Topic: "order-events",
Key: []byte("user-123"), // Key = 用户ID
Value: []byte(`{"action":"pay"}`),
}

// 不设置 Key(默认 nil)→ Sticky Partitioner 分散路由
record := &kgo.Record{
Topic: "order-events",
Value: []byte(`{"action":"pay"}`),
}

路由规则:

if Key != nil:
partition = hash(Key) % numPartitions // 同 Key 永远同 Partition
else:
Sticky Partitioner: 攒满一批发同一个 Partition,下一批随机换
Key 的选择 效果 适用场景
userID 同一用户的消息有序 用户行为流、状态变更
orderID 同一订单的消息有序 订单生命周期事件
nil(不设) 分散到所有 Partition,吞吐最大化 日志采集、无序指标

Key 不会让集群退化成单点

常见疑问:"指定 Key 后消息都路由到同一个 Partition,那不就只打一台 Broker 了?"

答案:不同的 Key 会被 hash 到不同的 Partition,分布在不同 Broker 上.假设你有百万用户,百万个不同的 Key 被均匀打散到 6 个 Partition → 3 台 Broker 均摊负载.Key 保证的是"同一个 Key 有序",不是"所有消息去一个地方".

hash("user-001") % 6 = 2 → P2 → Broker 2
hash("user-002") % 6 = 5 → P5 → Broker 2
hash("user-003") % 6 = 0 → P0 → Broker 0
hash("user-004") % 6 = 3 → P3 → Broker 0
hash("user-005") % 6 = 1 → P1 → Broker 1
hash("user-006") % 6 = 4 → P4 → Broker 1

只有极端场景(90% 流量来自同一个 Key)才会导致热点,这是"数据倾斜"问题而非 Key 机制的缺陷.

Partition 内有序,跨 Partition 无序

理解了 Key 路由后,再看顺序性就很清晰了.

假设你的服务按时间顺序产生了 5 条消息,它们因为 Key 不同被路由到了不同 Partition:

生产顺序: msg-1(user-A), msg-2(user-B), msg-3(user-A), msg-4(user-B), msg-5(user-A)

按 Key hash 路由后:
Partition 0: [msg-1] [msg-3] [msg-5] ← user-A 的三条,严格有序
Partition 1: [msg-2] [msg-4] ← user-B 的两条,严格有序

Kafka 的顺序保证仅限于单个 Partition 内部:

  • Partition 0 内,msg-1 一定在 msg-3 前面,msg-3 一定在 msg-5 前面——这个顺序不可逆,Consumer 读到的顺序永远如此
  • 但 msg-1(user-A)和 msg-2(user-B)谁先被消费?没有保证,因为它们在不同 Partition

这正是 Kafka 的设计意图:用 Partition 级有序换取集群级并行.业务上真正需要有序的消息(同一用户/同一订单)通过 Key 绑定到同一 Partition,天然有序;不相关的消息分散到不同 Partition,获得并行度.

顺序性陷阱

如果你不设 Key(Key=nil),消息会被 Sticky Partitioner 分散到不同 Partition——即使是同一用户的连续事件也可能被打散到不同 Partition,失去顺序保证.
需要顺序的场景,必须显式设置 Key.

分区数怎么定

不需要现在精确计算,但记住经验法则:

分区数 ≥ max(期望的 Producer 并发数, Consumer Group 中 Consumer 实例数)

分区创建后可以增加但不能减少(减少会导致数据丢失),所以初始别设太少.典型生产环境:6-12 个分区/topic 起步.


Offset:分区内消息的唯一地址

每条消息在其 Partition 内有一个单调递增的 64 位整数编号,叫 Offset.它是消息的"身份证号":

Partition 0:
offset 0 → {"user":"alice","action":"login"}
offset 1 → {"user":"bob","action":"purchase"}
offset 2 → {"user":"alice","action":"logout"}
offset 3 → (下一条消息将写在这里)

Offset 的关键特性

特性 说明
分区级唯一 Partition 0 的 offset 5 和 Partition 1 的 offset 5 是不同消息
不可变 一旦分配不会改变,消息也不会修改
不连续也合法 压缩后(log compaction)可能有空洞
Consumer 自己管理 Broker 不推送,Consumer 用 offset 告诉 Broker"我要从哪里开始读"

与其他系统的对比

系统 消费位置管理 类比
Redis Stream 每个 consumer group 有 last-delivered-id 类似,但 Redis 由服务端追踪
RabbitMQ Broker 记录 ack 状态,消费后删除 完全不同:推模型,消息即阅即焚
Kafka Consumer 提交 offset,消息保留到过期 拉模型,消息可重复消费
Go channel 读走就没了,无法回溯 Kafka 的 offset 允许"倒带"

Offset 带来的超能力

因为消息持久保留且有 offset 寻址,你可以:

  • 重新消费:把 offset 重置到过去某个时间点,重跑全部数据
  • 多个消费者独立消费:每个 group 各自维护 offset,互不影响
  • 故障恢复:crash 后从上次提交的 offset 续读,不丢不漏

全局视角:这些概念的关系

            Kafka Cluster
│
┌──────────┼──────────┐
▼ ▼ ▼
Broker 0 Broker 1 Broker 2
│ │ │
▼ ▼ ▼
┌─────────────────────────────────┐
│ Topic: order-events │
│ │
│ Partition 0 (leader: Broker 0) │
│ offset 0, 1, 2, 3 ... │
│ │
│ Partition 1 (leader: Broker 1) │
│ offset 0, 1, 2, 3 ... │
│ │
│ Partition 2 (leader: Broker 2) │
│ offset 0, 1, 2, 3 ... │
└─────────────────────────────────┘

一句话总结它们的层级关系:

Cluster 包含多个 Broker;一个 Topic 切分为多个 Partition;每个 Partition 分布在某个 Broker 上;Partition 内的每条消息有唯一 Offset.


动手验证:用 CLI 看到这些概念

后续会搭完整环境,这里先给个预览——用 kafka-topics.sh 可以直观看到 Topic 的分区分布:

# 创建一个 3 分区、2 副本的 topic
kafka-topics.sh --create \
--topic order-events \
--partitions 3 \
--replication-factor 2 \
--bootstrap-server localhost:9092

# 查看 topic 详情
kafka-topics.sh --describe --topic order-events --bootstrap-server localhost:9092

输出示例:

Topic: order-events  PartitionCount: 3  ReplicationFactor: 2
Partition: 0 Leader: 1 Replicas: 1,2 Isr: 1,2
Partition: 1 Leader: 2 Replicas: 2,0 Isr: 2,0
Partition: 2 Leader: 0 Replicas: 0,1 Isr: 0,1

对照着看:

  • 3 个 Partition,分别由 Broker 1、2、0 当 Leader
  • 每个 Partition 有 2 个副本(Replicas),分布在不同 Broker 上
  • ISR(In-Sync Replicas)表示当前与 Leader 同步的副本集合

小结

概念 一句话 类比
Broker 集群中的一台机器/进程 K8s 中的 Node
Topic 消息的逻辑分类 数据库中的表名
Partition Topic 的物理切片,有序日志 分库分表中的分片
Offset 消息在 Partition 内的序号 数组下标,但只增不减

下一篇我们看 Producer 端:消息写入时怎么选 Partition、acks 机制如何保证写入可靠性、批量发送的原理.