Skip to content

Kafka分区顺序与Offset

Kafka 的分区和 offset 是最容易被低估的知识点。很多线上问题不是“Kafka 不可靠”,而是开发者没有理解分区顺序、消费组分配和 offset 提交。

分区解决什么问题

一个 Topic 如果只有一个 Partition,就像只有一条收银通道,所有消息排队写入和读取。分区的作用是把一个 Topic 拆成多条并行通道。

mermaid
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:由分区器选择,通常用于均衡分散。
mermaid
flowchart TD
    A["Producer 发送消息"] --> B{"是否指定 Partition?"}
    B -->|"是"| C["写入指定 Partition"]
    B -->|"否"| D{"是否有 Key?"}
    D -->|"是"| E["按 Key 哈希选择 Partition"]
    D -->|"否"| F["按分区器策略分配"]

如果业务要求同一个订单的事件顺序处理,就应该用订单 ID 作为 key。

java
kafkaTemplate.send("order.events", String.valueOf(orderId), event);

这样同一个 orderId 的消息会进入同一个分区,消费者读取时就是分区内顺序。

顺序消息的边界

假设一个订单有三个事件:

  1. ORDER_CREATED
  2. ORDER_PAID
  3. ORDER_CANCELED

如果三条消息 key 都是 orderId=1001,它们会进入同一个分区,Kafka 可以保证消费顺序。

如果 key 不固定,三条消息可能分散到不同分区:

mermaid
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。
mermaid
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。

消费者数量和分区数量

在同一个消费组内,一个分区同一时刻只能分配给一个消费者;一个消费者可以处理多个分区。

分区数消费者数结果
311 个消费者处理 3 个分区
33每个消费者处理 1 个分区
353 个消费者工作,2 个消费者空闲

所以盲目增加消费者不一定提升吞吐。如果分区数只有 3,消费者扩到 20 也最多 3 个在处理。

Offset 是什么

Offset 是 Partition 内的消息位置。每个消费组都会维护自己的 offset。

mermaid
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:“这个消费组已经处理到这里了”。

自动提交

自动提交简单,但有风险:

mermaid
flowchart TD
    A["拉取消息"] --> B["自动提交 Offset"]
    B --> C["处理业务"]
    C --> D{"业务成功?"}
    D -->|"否"| E["消息不会再从旧 Offset 自动重放"]

如果 offset 已提交,但业务处理失败或服务宕机,就可能漏处理。

手动提交

手动提交更适合生产业务:

mermaid
flowchart TD
    A["拉取消息"] --> B["处理业务"]
    B --> C{"业务成功?"}
    C -->|"是"| D["提交 Offset"]
    C -->|"否"| E["不提交,等待重试或进入异常处理"]

手动提交不能消灭重复消费,但可以避免“业务没成功,进度已经前进”的漏消费问题。

Rebalance

Rebalance 是消费组重新分配分区的过程。常见触发条件:

  • 消费者上线。
  • 消费者下线。
  • 消费者长时间没有心跳。
  • Topic 分区数量变化。
mermaid
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 的很多线上问题,本质上都是这三个概念没有设计好。