Kafka分区顺序与Offset
Kafka 的分区和 offset 是最容易被低估的知识点。很多线上问题不是“Kafka 不可靠”,而是开发者没有理解分区顺序、消费组分配和 offset 提交。
分区解决什么问题
一个 Topic 如果只有一个 Partition,就像只有一条收银通道,所有消息排队写入和读取。分区的作用是把一个 Topic 拆成多条并行通道。
flowchart TD
A["Topic: user.behavior"] --> B["Partition 0"]
A --> C["Partition 1"]
A --> D["Partition 2"]
B --> E["Consumer A"]
C --> F["Consumer B"]
D --> G["Consumer C"]分区带来两个能力:
- 写入并行:不同分区可以分散到不同 Broker。
- 消费并行:同一个消费组内多个消费者可以分摊不同分区。
代价是:Kafka 只能保证单个 Partition 内消息顺序,不能保证 Topic 全局顺序。
消息如何选择分区
常见规则:
- 指定 partition:直接写入指定分区。
- 指定 key:根据 key 哈希选择分区。
- 不指定 key:由分区器选择,通常用于均衡分散。
flowchart TD
A["Producer 发送消息"] --> B{"是否指定 Partition?"}
B -->|"是"| C["写入指定 Partition"]
B -->|"否"| D{"是否有 Key?"}
D -->|"是"| E["按 Key 哈希选择 Partition"]
D -->|"否"| F["按分区器策略分配"]如果业务要求同一个订单的事件顺序处理,就应该用订单 ID 作为 key。
kafkaTemplate.send("order.events", String.valueOf(orderId), event);这样同一个 orderId 的消息会进入同一个分区,消费者读取时就是分区内顺序。
顺序消息的边界
假设一个订单有三个事件:
ORDER_CREATEDORDER_PAIDORDER_CANCELED
如果三条消息 key 都是 orderId=1001,它们会进入同一个分区,Kafka 可以保证消费顺序。
如果 key 不固定,三条消息可能分散到不同分区:
flowchart TD
A["ORDER_CREATED"] --> B["Partition 0"]
C["ORDER_PAID"] --> D["Partition 1"]
E["ORDER_CANCELED"] --> F["Partition 2"]
B --> G["Consumer A"]
D --> H["Consumer B"]
F --> I["Consumer C"]这时三个消费者并行处理,业务上就可能出现“先取消、后支付、再创建”的错乱。
Consumer Group
Consumer Group 是 Kafka 消费模型的核心。
- 同一个 group 内:多个消费者分摊同一个 Topic 的分区,一条消息只会被组内一个消费者处理。
- 不同 group 之间:互不影响,每个 group 都能完整消费 Topic。
flowchart TD
A["Topic: order.created"] --> B["Partition 0"]
A --> C["Partition 1"]
A --> D["Partition 2"]
B --> E["Group A Consumer 1"]
C --> F["Group A Consumer 2"]
D --> F
B --> G["Group B Consumer 1"]
C --> G
D --> H["Group B Consumer 2"]这解释了一个常见问题:为什么两个服务监听同一个 Topic,只有一个服务收到消息?原因通常是它们用了同一个 groupId。如果你希望两个服务都收到消息,它们必须使用不同的 groupId。
消费者数量和分区数量
在同一个消费组内,一个分区同一时刻只能分配给一个消费者;一个消费者可以处理多个分区。
| 分区数 | 消费者数 | 结果 |
|---|---|---|
| 3 | 1 | 1 个消费者处理 3 个分区 |
| 3 | 3 | 每个消费者处理 1 个分区 |
| 3 | 5 | 3 个消费者工作,2 个消费者空闲 |
所以盲目增加消费者不一定提升吞吐。如果分区数只有 3,消费者扩到 20 也最多 3 个在处理。
Offset 是什么
Offset 是 Partition 内的消息位置。每个消费组都会维护自己的 offset。
flowchart TD
A["Partition 0"] --> B["Offset 0"]
B --> C["Offset 1"]
C --> D["Offset 2"]
D --> E["Offset 3"]
F["Group A 已提交 Offset 3"] --> E
G["Group B 已提交 Offset 1"] --> C注意:offset 是“消费组维度”的,不是 Topic 全局维度。不同消费组可以读到不同位置,这就是 Kafka 支持多系统独立消费和回放的基础。
自动提交和手动提交
Kafka 消费后需要提交 offset。提交 offset 的意思是告诉 Kafka:“这个消费组已经处理到这里了”。
自动提交
自动提交简单,但有风险:
flowchart TD
A["拉取消息"] --> B["自动提交 Offset"]
B --> C["处理业务"]
C --> D{"业务成功?"}
D -->|"否"| E["消息不会再从旧 Offset 自动重放"]如果 offset 已提交,但业务处理失败或服务宕机,就可能漏处理。
手动提交
手动提交更适合生产业务:
flowchart TD
A["拉取消息"] --> B["处理业务"]
B --> C{"业务成功?"}
C -->|"是"| D["提交 Offset"]
C -->|"否"| E["不提交,等待重试或进入异常处理"]手动提交不能消灭重复消费,但可以避免“业务没成功,进度已经前进”的漏消费问题。
Rebalance
Rebalance 是消费组重新分配分区的过程。常见触发条件:
- 消费者上线。
- 消费者下线。
- 消费者长时间没有心跳。
- Topic 分区数量变化。
flowchart TD
A["Group 内消费者变化"] --> B["触发 Rebalance"]
B --> C["暂停部分消费"]
C --> D["重新分配 Partition"]
D --> E["消费者从已提交 Offset 继续消费"]Rebalance 期间如果 offset 提交不合理,就容易重复消费或短暂消费停顿。消费者处理时间过长也可能被认为失联,从而触发 Rebalance。
分区数怎么设置
分区数没有万能公式,需要结合吞吐、消费者并行度、顺序要求和未来扩展。
可以按下面思路估算:
- 需要多少消费并行度:分区数至少要大于等于目标消费者并发数。
- 单分区吞吐是否够:如果单分区写入或读取成为瓶颈,需要增加分区。
- 是否有顺序要求:同一业务 key 必须进入同一分区,分区过多不破坏局部顺序,但会增加管理成本。
- 是否可能扩容:分区数后续可以增加,但增加后 key 到分区的映射可能变化,历史顺序要特别注意。
生产中不要随意把分区数设得特别大。分区太多会增加文件句柄、内存、Controller 管理成本、Rebalance 成本。
常见问题
Kafka 能保证顺序吗
能,但只保证同一个 Partition 内顺序。要让同一订单、同一用户、同一设备的消息有序,需要使用稳定 key。
增加分区会影响顺序吗
可能会。增加分区后,key 哈希到哪个分区的结果可能变化。对于严格顺序业务,需要谨慎扩分区,或者使用业务路由策略保证同一 key 的新旧消息不乱。
消费很慢怎么办
先判断瓶颈:
- 单条业务处理慢:优化业务逻辑或异步化。
- 分区数太少:增加分区和消费者。
- 下游数据库慢:限流、批处理、缓存、削峰。
- Rebalance 频繁:检查消费者心跳和处理时间。
小结
Partition 决定吞吐和顺序边界,Consumer Group 决定消息是竞争消费还是广播消费,Offset 决定消费进度和故障恢复。Kafka 的很多线上问题,本质上都是这三个概念没有设计好。
