一句话总结
Kafka 的数据模型是一个只追加的分布式日志:消息写进 Topic 的某个 Partition,每条消息有唯一递增的 Offset;多个 Broker 组成集群对外提供读写.
把握住"append-only log + offset 寻址"这个心智模型,后续所有机制都是它的衍生.
用一条消息的旅程串起全部概念
假设你的 Go 服务要把一条订单事件发到 Kafka,消费端读取后写入数据仓库.整个流程涉及的核心角色:
|
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.
|
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,实现负载均衡:
|
关键结论:Partition 是负载分散的单位.6 个 Partition 分布在 3 个 Broker → 读写压力被 3 台机器分摊.
消息路由:Key 决定去哪个 Partition
Producer 发消息时可以附带一个 Key(任意 []byte),Kafka 根据 Key 决定消息写入哪个 Partition:
|
路由规则:
|
| Key 的选择 | 效果 | 适用场景 |
|---|---|---|
userID |
同一用户的消息有序 | 用户行为流、状态变更 |
orderID |
同一订单的消息有序 | 订单生命周期事件 |
nil(不设) |
分散到所有 Partition,吞吐最大化 | 日志采集、无序指标 |
Key 不会让集群退化成单点
常见疑问:"指定 Key 后消息都路由到同一个 Partition,那不就只打一台 Broker 了?"
答案:不同的 Key 会被 hash 到不同的 Partition,分布在不同 Broker 上.假设你有百万用户,百万个不同的 Key 被均匀打散到 6 个 Partition → 3 台 Broker 均摊负载.Key 保证的是"同一个 Key 有序",不是"所有消息去一个地方".
|
只有极端场景(90% 流量来自同一个 Key)才会导致热点,这是"数据倾斜"问题而非 Key 机制的缺陷.
Partition 内有序,跨 Partition 无序
理解了 Key 路由后,再看顺序性就很清晰了.
假设你的服务按时间顺序产生了 5 条消息,它们因为 Key 不同被路由到了不同 Partition:
|
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.
分区数怎么定
不需要现在精确计算,但记住经验法则:
|
分区创建后可以增加但不能减少(减少会导致数据丢失),所以初始别设太少.典型生产环境:6-12 个分区/topic 起步.
Offset:分区内消息的唯一地址
每条消息在其 Partition 内有一个单调递增的 64 位整数编号,叫 Offset.它是消息的"身份证号":
|
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 续读,不丢不漏
全局视角:这些概念的关系
|
一句话总结它们的层级关系:
Cluster 包含多个 Broker;一个 Topic 切分为多个 Partition;每个 Partition 分布在某个 Broker 上;Partition 内的每条消息有唯一 Offset.
动手验证:用 CLI 看到这些概念
后续会搭完整环境,这里先给个预览——用 kafka-topics.sh 可以直观看到 Topic 的分区分布:
|
输出示例:
|
对照着看:
- 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 机制如何保证写入可靠性、批量发送的原理.