在实时对话 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 的闲置回收周期要调好,太短会频繁创建销毁,太长会占内存;背压必须设计进每层,否则突发流量会打爆内存;死信队列要有人工补偿机制,否则永久错误消息会把会话永久卡住。核心原则只有一句话:把会话状态的序列化作用域收拢到单点,用结构保序,而不是靠事后协调。

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

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

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

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

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

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

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

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

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

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

是连续水位线提交机制:只提交最高连续完成的偏移,空隙处阻塞提交。
这套方案的最终启发是:在分布式系统里,“顺序”不是一个可以事后补偿的属性,而是一种需要从架构初始就刻进数据路径的结构约束。当你用 Kafka 做管道时,记住分区顺序只是起点,真正的会话有序得靠应用层自己负责任——用一致性哈希锁定亲和性,用单 goroutine 串行化关键路径,用连续水位线保证可恢复性,你就能在几十万 TPS 的量级上仍然睡得着觉。
内容与图片版权归原作者所有 · 原文: https://www.infoq.com/articles/apache-kafka-golang-session-ordered-pipeline/?utm_campaign=infoq_content&utm_source=infoq&utm_medium=feed&utm_term=global