在实时对话 AI 场景里,消息顺序本身就是语义。用户先发“帮我订一张去伦敦的机票”,紧接着又补一句“算了,改成巴黎”,如果第二条消息先被处理,AI 就会回一个伦敦确认单而不是巴黎改签。这个错误在单个会话里看是小概率事件,但当成千上万个会话同时流过一条多级管道时,哪怕极低的乱序率都会毁掉产品体验。我们为一家服务企业客户的实时对话 AI 平台设计了核心 Kafka 消息层,用 Go 实现,要求严格:同一会话内的消息必须按发送顺序依次经过 ingest(数据接入)、NLU(自然语言理解)、LLM(大模型推理)、delivery(交付)四个阶段,不管某个阶段耗时多长。

难点在于,Kafka 只保证分区内的顺序,但一个分区会承载大量会话的交错流量。要做到跨会话并行、会话内有序,就得在应用层自己加路由逻辑。我们第一版用了一个固定的扁平 worker 池,把会话哈希到 worker 上。在持续负载下,一个正在重试的会话会把同 worker 上无关的会话也堵住,而重试协调又引入了不断增大的内存队列,最终被换成了两级 worker 层级:第一层是轻量级 dispatch worker(分发协程),从 Kafka 拉记录,按 session ID 做一致性哈希(一种把相同 key 稳定映射到同一节点的哈希方法),保证同一会话的消息总落到同一个 dispatch worker;第二层是每个活跃会话独立的 goroutine(Go 的轻量级并发执行单元),它按到达顺序逐条处理消息,会话创建按需,超时无活动就销毁,所以资源占用和活跃会话数成正比。每个会话 goroutine 同一时间只处理一条消息,跨会话的并行度天然无上限,会话内的串行则由结构保证。

为什么不干脆用 Kafka 分区来做排序?分区对应会话在几千并发会话下不可行,分区数量会超过运维上限并拖慢 broker(Kafka 的服务器节点)。单线程处理每个分区也不行,因为分区里交杂着多个会话,会话级排序还是需要应用层协调。一致性哈希既保证了严格会话内有序,又给了高跨会话并行,还不会让分区数无限膨胀。

重试是破坏顺序的典型场景。第一版用独立的 retry goroutine 排队等待消息,状态机复杂且内存不可控。现在改成在会话 goroutine 内部就地重试:处理一条消息,成功就标记完成并继续下一条;遇到瞬时故障(比如下游服务超时),goroutine 进入指数退避(一种不断拉长等待时间的重试策略)加抖动(随机化退避间隔防止惊群)后重试同一条消息。因为 goroutine 在重试期间是阻塞的,该会话后续消息只能堆积在通道里,不可能抢跑失败的消息。假如是永久性错误(比如消息数据非法),就直接踢到死信队列(DLQ,专门存放无法处理的消息的队列),不重试。这样省掉了外部重试 goroutine,没有额外的队列、通道和状态重置,消息 4 在消息 3 完成前物理上无法执行。

多会话 goroutine 乱序完成时,如果朴素提交最新完成的偏移量(offset,每条消息在分区中的序号),会出大问题:假设偏移 104 先完成而 102 还没完成,提交 104 后一旦崩溃重启,102 就被永久跳过。我们的解法是维护每个分区的“连续水位线”(contiguous watermark):一个跟踪组件分别记录 InFlight(正在处理的偏移)和 Completed(已完成的偏移),定期提交时只推进到连续完成的最高偏移,中间的空隙(没有完成的消息)就是屏障,阻止不安全推进。崩溃后从水位线处重放,可能重放部分已完成消息,因此我们给每条消息带一个稳定的事件 ID 做去重,把重复影响降到最小,但外部副作用还要依赖下游服务的幂等性,所以不是精确一次语义。这是刻意的取舍:宁愿重放少量消息也不冒静默丢失数据的风险。

生产级有序保证还离不开运维层面的加固:再均衡时排空、背压(下游慢时限制上游流速的机制)、待处理记录缓冲、卡死偏移检测、幂等、优雅停机、死信队列处理,这些在真实故障下缺一不可。这套管道在线上已处理超过 4000 万条消息,在数千并发会话下没有观测到一次乱序;在测试中发送侧吞吐已突破每秒 10 万条消息,零发送错误。最让我意外的是,这个有序性不是靠压着吞吐换来的,而是由一致性哈希加每会话单 goroutine 的结构性设计天然决定的。

如果你也在做类似的消息管道——不一定是 AI,只要是按业务 key(用户、订单、设备)保序且要跨 key 并行的场景,比如游戏事件流、金融交易流水、IoT 指令下发,这套思路可以直接套用:第一层用一致性哈希做亲和路由,第二层每个 key 一个轻量 worker 顺序消费,重试就地执行而不另开线程,偏移提交只在连续完成处前进。要注意的坑有:会话 goroutine 的闲置回收周期要调好,太短会频繁创建销毁,太长会占内存;背压必须设计进每层,否则突发流量会打爆内存;死信队列要有人工补偿机制,否则永久错误消息会把会话永久卡住。核心原则只有一句话:把会话状态的序列化作用域收拢到单点,用结构保序,而不是靠事后协调。

Accelerating Performance by Incrementally Integrating Rust into Existing Codebas
Accelerating Performance by Incrementally Integrating Rust into Existing Codebas

展示的是把 Rust 增量融入现有代码库以加速性能(原文配图,此处保留占位)。

Building a Session-Ordered Kafka Pipeline in Go
Building a Session-Ordered Kafka Pipeline in Go

是整体管道的会话有序数据流。

Building Reusable Evaluation Frameworks for Agentic AI Products
Building Reusable Evaluation Frameworks for Agentic AI Products

给出了两层 worker 层级中一致性哈希的分布示意。

The Reinvention of the Dev Team
The Reinvention of the Dev Team

是会话 worker 与每条会话 goroutine 的协作关系。

Building Resilient Platforms: Insights from 20+ Years in Mission-Critical Infras
Building Resilient Platforms: Insights from 20+ Years in Mission-Critical Infras

是分区内交错会话流量如何被路由到不同 worker。

对比了分区内排序与一致性哈希路由的差异。

是会话 goroutine 就地重试的时序。

展示了阻塞式重试如何让后续消息无法抢跑。

列出了两层错误分类:透明错误走退避重试,毒丸消息直接进 DLQ。

画出了乱序完成偏移的危险场景。

是连续水位线提交机制:只提交最高连续完成的偏移,空隙处阻塞提交。

这套方案的最终启发是:在分布式系统里,“顺序”不是一个可以事后补偿的属性,而是一种需要从架构初始就刻进数据路径的结构约束。当你用 Kafka 做管道时,记住分区顺序只是起点,真正的会话有序得靠应用层自己负责任——用一致性哈希锁定亲和性,用单 goroutine 串行化关键路径,用连续水位线保证可恢复性,你就能在几十万 TPS 的量级上仍然睡得着觉。

阅读原文 → 返回 AI 技术文档

内容与图片版权归原作者所有 · 原文: https://www.infoq.com/articles/apache-kafka-golang-session-ordered-pipeline/?utm_campaign=infoq_content&utm_source=infoq&utm_medium=feed&utm_term=global