Kafka 把 Topic 拆成多个 Partition。每个 Partition 是有序、可追加的日志;消费者通过 Offset 记录读到哪里。Kafka 的并行度、顺序和扩容都围绕 Partition 展开。
核心对象
| 对象 | 作用 |
|---|---|
| Broker | 保存分区并处理生产、消费请求 |
| Topic | 一类事件的逻辑名称 |
| Partition | 有序追加日志,也是并行单位 |
| Replica | 分区副本,用于故障恢复 |
| Consumer Group | 协作消费同一 Topic 的消费者集合 |
| Offset | 消费者在某个分区中的位置 |
| Lag | 最新 Offset 与已提交 Offset 的差值 |
同一消费组中,一个 Partition 在同一时刻只交给一个消费者。若有 6 个 Partition、3 个消费者,每个消费者大致负责 2 个;扩到 8 个消费者时,多出的 2 个通常没有分区可处理。
Kafka 如何存储
1 | Producer -> 选择 Partition -> 写 Leader 日志 |
消息消费后不会立即删除。Kafka 按保留时间、容量或日志压缩策略管理数据,因此消费者可以调整 Offset 后重放历史事件。
ZooKeeper 与 KRaft
ZooKeeper 是独立的分布式协调系统,早期 Kafka 用它保存集群元数据和协助控制器选举。现代 Kafka 已把元数据管理迁移到内置的 KRaft 共识模式,新集群不应再把 ZooKeeper 当作 Kafka 消息传输链路中的必备组件。
1 | 旧架构:Kafka Broker + ZooKeeper 管理元数据 |
ZooKeeper 并不保存 Kafka 的业务消息;消息仍在 Kafka Partition 日志中。维护存量集群时要先确认版本和运行模式,再按官方迁移文档操作。
怎样保证设备内顺序
生产者使用稳定的 device_id 作为消息 Key,让同一设备进入同一 Partition:
1 | device-001 seq=101 -> Partition 3 |
Kafka 只保证单个 Partition 内的读取顺序,不保证整个 Topic 的全局顺序。端到端仍可能因生产重试、消费者失败和数据库写入出现重复或乱序,所以事件应包含 event_id、设备序号和设备时间。
Offset 何时提交
1 | 拉取消息 -> 处理业务 -> 持久化成功 -> 提交 Offset |
先提交再处理可能丢业务结果;处理成功后提交可能在崩溃时重复消费。因此常见目标是“至少一次投递 + 幂等处理”,而不是口头承诺所有外部系统都恰好一次。
Lag 与扩容
Kafka 不会直接创建 Worker。监控系统或 KEDA 读取消费组 Lag,根据规则修改 Kubernetes Deployment 副本数:
1 | Kafka Lag 上升 |
扩容前要确认 Partition 数量和数据库容量。只增加 Pod 而不增加可用 Partition,不会继续提升并行度;消费者过多还可能让 Rebalance 更频繁。
一条 IoT 上行链路
1 | 设备 -> MQTT Broker -> 规则/桥接 -> Kafka |
MQTT 负责设备连接和 Topic 路由,Kafka 负责云端缓冲、重放与多下游消费,两者不是替代关系。