跳过导航

消息顺序性到底应该如何保证?从 Kafka 分区到业务状态机

约 6 分钟...次浏览
专栏分布式系统与消息队列第 4 篇

业务常说“消息必须有序”,但很少先定义有序的范围。是所有订单全局有序,还是同一个订单的创建、支付、取消有序?定义不同,系统吞吐和复杂度可能相差几个数量级。

一、三种顺序语义

全局有序

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 路由同一分区是基础,用业务版本、幂等和状态机应对重试、并发、扩分区与迟到才是完整方案。全局严格有序代价极高,除非业务确实需要,否则应把顺序约束缩小到最小业务范围。

分享:
文章作者:狼码纪
版权声明:本博客所有文章除特别声明外,均采用 CC BY-NC-SA 4.0 许可协议。文章可能参考了其他优秀文章,如有侵权请联系删除。