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
2
3
4
5
Producer -> 选择 Partition -> 写 Leader 日志
-> Followers 复制

Partition 0: offset 0, 1, 2, 3 ...
Partition 1: offset 0, 1, 2, 3 ...

消息消费后不会立即删除。Kafka 按保留时间、容量或日志压缩策略管理数据,因此消费者可以调整 Offset 后重放历史事件。

ZooKeeper 与 KRaft

ZooKeeper 是独立的分布式协调系统,早期 Kafka 用它保存集群元数据和协助控制器选举。现代 Kafka 已把元数据管理迁移到内置的 KRaft 共识模式,新集群不应再把 ZooKeeper 当作 Kafka 消息传输链路中的必备组件。

1
2
旧架构:Kafka Broker + ZooKeeper 管理元数据
新架构:Kafka Broker/Controller + KRaft 元数据日志

ZooKeeper 并不保存 Kafka 的业务消息;消息仍在 Kafka Partition 日志中。维护存量集群时要先确认版本和运行模式,再按官方迁移文档操作。

怎样保证设备内顺序

生产者使用稳定的 device_id 作为消息 Key,让同一设备进入同一 Partition:

1
2
3
device-001 seq=101 -> Partition 3
device-001 seq=102 -> Partition 3
device-002 seq=77 -> Partition 1

Kafka 只保证单个 Partition 内的读取顺序,不保证整个 Topic 的全局顺序。端到端仍可能因生产重试、消费者失败和数据库写入出现重复或乱序,所以事件应包含 event_id、设备序号和设备时间。

Offset 何时提交

1
拉取消息 -> 处理业务 -> 持久化成功 -> 提交 Offset

先提交再处理可能丢业务结果;处理成功后提交可能在崩溃时重复消费。因此常见目标是“至少一次投递 + 幂等处理”,而不是口头承诺所有外部系统都恰好一次。

Lag 与扩容

Kafka 不会直接创建 Worker。监控系统或 KEDA 读取消费组 Lag,根据规则修改 Kubernetes Deployment 副本数:

1
2
3
4
5
6
7
Kafka Lag 上升
-> KEDA Scaler 采集指标
-> 修改目标副本数
-> Kubernetes 创建 Pod
-> 新消费者加入 Group
-> Group Rebalance
-> Partition 重新分配

扩容前要确认 Partition 数量和数据库容量。只增加 Pod 而不增加可用 Partition,不会继续提升并行度;消费者过多还可能让 Rebalance 更频繁。

一条 IoT 上行链路

1
2
3
设备 -> MQTT Broker -> 规则/桥接 -> Kafka
-> Worker 校验、去重、排序 -> 时序数据库
-> 告警、报表和实时看板

MQTT 负责设备连接和 Topic 路由,Kafka 负责云端缓冲、重放与多下游消费,两者不是替代关系。

延伸阅读

站内搜索

没有找到内容!