当数千个会话共享同一个分区(partition)时,仅依赖 Kafka 的分区内顺序是不够的。在聊天系统和顺序任务中,一条修正消息,例如“将目的地设为巴黎”,可能在“预订前往伦敦的航班”之后到达,但如果它先于第一条消息执行,就会破坏上下文。Joshua Oluikpe 在 InfoQ 中介绍的解决方案,是将部分排序责任转移到使用 Go 编写的应用层。
会话内排序,会话间并行
该设计依赖两级工作者。分发层接收 Kafka 记录,并使用一致性哈希按照会话 ID进行路由,使同一会话的消息到达同一个工作者,同时允许不同会话分配给多个工作者。随后,每个活跃会话都会获得一个专用的 goroutine,一次处理一条消息,并按照到达顺序处理。
这样,运行缓慢的会话不会像扁平工作者池那样阻碍其他会话;在扁平工作者池中,单个会话的重试可能会完全占用一个工作者。goroutine 按需创建,并在一段时间不活跃后移除,因此资源消耗取决于活跃会话数量,而不是可能存在的会话总数。
在不越过消息的情况下重试
重试在同一个会话的 goroutine 内进行。发生临时错误时,例如下游服务超时,工作者会使用带有一定随机性的指数退避等待,然后重新处理同一条消息,再继续处理后续消息。由于后续消息会留在发生故障的消息之后的通道中,因此第四条消息不可能越过第三条消息。
而对于不可修复的错误,例如无效数据,则会直接进入死信队列(DLQ)。这种对临时错误和最终错误的区分,消除了使用独立状态机管理重试的必要,但也意味着某个会话的消息可能会在退避期间不断累积。
安全提交与故障恢复
当多个会话并行处理各自的消息时,不能只提交已完成的最高 offset。如果 offset 104 在 102 之前完成,并随后提交 104,那么消费者发生故障时可能会永久跳过 102。因此,系统会针对每个分区跟踪正在处理的 offset 和已完成的 offset,并且只将提交点推进到最高的、连续完成的 offset。
因此,重启可能会从更早的 offset 开始,并重新处理部分消息。该设计使用固定的事件 ID 来帮助去重,但并不声称实现“恰好一次”执行;外部影响仍取决于下游服务以幂等方式处理重复消息的能力。在极少数情况下,在指定超时后,可以跳过一个长期卡住的间隙,以恢复分区进度,同时将其记录为运行故障并把消息导向 DLQ。这是在持续运行能力与处理完整性之间进行的明确权衡。
真正投入生产还需要什么?
仅凭核心逻辑无法完成排序保证。在分区重新分配期间,消费者会停止接收新记录,等待当前工作完成,然后提交最后一个连续 watermark,再转移所有权。当工作者通道已满时,系统会暂时停止拉取消息,而不是丢弃消息;对于已经预先拉取的记录,则存储在各分区专用的缓冲区中,以免更新的消息超过它们。
系统还需要逐个检测卡住的 offset,跨越异步的 Kafka 边界关联跟踪信息,在停止拉取期间保留自动扩展所需的滞后信号,并在关闭前排空工作、写入 DLQ。产品还会在发送端应用一致性哈希,以确保异步生产时的消息顺序。
结果与限制
在一项包含 50,000 条消息、10 个会话,并向 10% 的消息注入临时错误的合成测试中,系统达到了每秒 14,027 条消息的处理速度,没有出现顺序违规,中位延迟为 11.7 毫秒。在包含 1,000 个会话的负载测试中,实际发送速率达到每秒 48,805 条消息,目标为每秒 50,000 条,没有发送错误;但尾部延迟超过了 200 毫秒的预算,原因是 API Gateway 的消息速率配额,而不是会话排序机制。
作者称,该系统已经处理了超过 4,000 万条生产消息,没有观察到明显的顺序违规,仅有 460 条消息被导向 DLQ,约为 0.002%。不过,这一数字基于运行监控,而不是对每条消息的完整顺序进行检查;此外,更高负载的测试也没有独立验证每个会话的顺序。因此,该设计的实际价值在于将排序从 Kafka 层面的假设转变为应用强制实施的属性,同时仍然需要明确的顺序测试,以及下游服务中的幂等性保证。