消息顺序性到底应该如何保证?从 Kafka 分区到业务状态机
业务常说“消息必须有序”,但很少先定义有序的范围。是所有订单全局有序,还是同一个订单的创建、支付、取消有序?定义不同,系统吞吐和复杂度可能相差几个数量级。
一、三种顺序语义
全局有序
Topic 中所有消息严格按唯一顺序消费。通常只能使用单分区和单消费执行链,扩展能力最弱。
分区有序
Kafka 保证单个 Partition 内记录按 offset 排列,不保证不同分区之间的相对顺序。
业务实体有序
只要求同一订单、账户或设备的事件有序。通常把聚合 ID 作为消息 key,使同一 key 稳定进入同一分区,是最实用的语义。
kafkaTemplate.send("order-events", orderId.toString(), event);
二、生产端也可能制造乱序
即使消息发往同一分区,多线程并发产生同一订单事件时,发送顺序可能本来就不确定。业务数据库提交顺序、Outbox 扫描顺序和实际发送顺序也可能不同。
应在事件中加入聚合版本:
{
"orderId": 1001,
"version": 7,
"type": "ORDER_PAID"
}
版本比时间戳可靠,因为不同机器时钟可能漂移,相同毫秒内也可能产生多条事件。
三、分区器与扩分区
同一 key 通常通过哈希映射到固定分区。但 Topic 增加分区后,hash(key) % partitionCount 的结果可能变化:旧事件在分区 2,新事件进入分区 7,两个分区被不同消费者并行处理,迁移窗口发生乱序。
扩分区前应评估:
- 是否能暂停生产并等待旧积压清空。
- 是否有自定义稳定映射或路由表。
- 消费端能否通过版本号容忍乱序。
- 是否更适合新建 Topic 并迁移。
Kafka 4.x 的新消费者再均衡协议可以改善扩缩容期间的停顿与分配过程,但不会把跨分区记录变成全局有序,也不会修复应用线程池造成的乱序。顺序保证仍以分区、key 路由和业务版本为边界。
四、消费者并发的隐藏破坏
Kafka 将一个分区交给组内一个消费者,但消费者内部若把消息提交到普通线程池:
for (var record : records) {
executor.submit(() -> handle(record));
}
任务执行完成顺序不再受 offset 约束。offset 提交也变得危险:offset 12 先完成并提交,而 offset 11 随后失败,11 可能被错误跳过。
需要并行时,可以按 key 做串行队列:同一 key 路由到同一工作槽,不同 key 并行。还要确保 offset 只推进到连续完成的最大位置。
五、失败重试如何破坏顺序
事件 v7 失败后被转到重试 Topic,主分区继续消费 v8。等 v7 重试成功时,状态已经到 v8。
严格顺序方案会暂停该分区,代价是一个坏消息阻塞其他 key。更细粒度方案会冻结单个聚合 key,但实现复杂。另一条路线是让状态机识别版本:
UPDATE order_projection
SET status = ?, version = ?
WHERE order_id = ? AND version = ?;
只有当前版本为 6 时才接受 v7。若 v8 先到,可暂存、重拉缺失事件或根据权威数据重建投影。
六、数据库提交与消息发送顺序
两个事务先后更新同一订单,Outbox 记录的自增 ID 不一定等同于业务版本,尤其在并发事务下,先获得 ID 的事务可能后提交。
可靠方案应在聚合更新时原子递增版本,并把该版本写入 Outbox。消费者依据业务版本,而不是 Kafka 到达时间推断因果顺序。
七、顺序与幂等必须一起设计
有序不代表不重复。消费者可能依次收到 v7、v7、v8。规则可以是:
incomingVersion == currentVersion + 1:正常处理。incomingVersion <= currentVersion:重复或迟到,幂等忽略。incomingVersion > currentVersion + 1:存在缺口,暂存并告警。
switch (Long.compare(event.version(), state.version() + 1)) {
case 0 -> apply(event);
case -1 -> log.debug("duplicate or stale event");
case 1 -> gapStore.save(event);
}
八、什么时候不必追求严格顺序
指标统计、日志采集、搜索索引刷新等场景常可使用最终状态覆盖,或按事件时间窗口聚合。为它们强行使用单分区,会牺牲吞吐与可用性。
先问业务真正不能接受的是“乱序到达”,还是“最终状态错误”。很多系统只需要后者不发生,通过版本化和幂等就足够。
九、典型方案选择
| 需求 | 推荐方案 |
|---|---|
| 所有消息严格有序 | 单分区、单执行链 |
| 同订单有序 | orderId 作为 key + 分区内串行 |
| 高吞吐且允许乱序到达 | 版本号 + 幂等状态机 |
| 重试仍需严格顺序 | 暂停分区或 key 级冻结 |
| 可重建读模型 | 保留事件日志,按版本重放 |
十、检查清单
- “有序”的范围是全局、分区还是单业务实体?
- 同一实体是否始终使用稳定且非空的 key?
- 生产端并发时是否生成单调业务版本?
- 扩分区是否会改变 key 映射?
- 消费者内部线程池是否打乱执行和提交顺序?
- 失败转重试 Topic 后如何处理后续版本?
- 是否同时处理重复、迟到与版本缺口?
- 能否从权威状态或事件日志重建投影?
总结
消息中间件能提供的通常是分区内顺序,而业务需要的是因果顺序。用聚合 ID 路由同一分区是基础,用业务版本、幂等和状态机应对重试、并发、扩分区与迟到才是完整方案。全局严格有序代价极高,除非业务确实需要,否则应把顺序约束缩小到最小业务范围。
相关文章
Exactly Once 是真实能力还是营销概念?先定义处理边界
区分投递一次、处理一次与业务效果一次,分析 Kafka 幂等生产者、事务、read-process-write 链路及外部数据库边界,说明 Exactly Once 的真实适用范围。
Kafka 消息为什么仍可能重复或丢失?生产、Broker、消费三段治理
沿生产者、Broker 副本和消费者三个环节分析 Kafka 消息丢失与重复的故障窗口,给出幂等生产、acks、ISR、手动提交和业务幂等方案。
Kafka 为什么能实现高吞吐?顺序写、批处理、页缓存与零拷贝
从分区日志、顺序追加、批量压缩、操作系统页缓存、sendfile、消费者拉取和横向扩展解释 Kafka 高吞吐的完整机制。