Skip to content

【Managed Agent】事件主干:为什么是Kafka

本文单独讨论Aviary的事件主干:为什么用户请求要先被持久化并转换为标准事件,Kafka在可靠投递、按对话保序、幂等消费和异步推进中的位置,以及这套设计如何从早期方案演进而来。

大纲

1. 问题从哪里来:一段对话为什么不能直接调用Agent

  • 长任务不能占住用户请求链路。
  • 请求已经确认成功后,进程重启、下游暂时故障或Worker扩缩容,都不应让这件事消失。
  • 同一段对话必须按顺序推进,不同对话又需要并行。
  • Agent任务天然是异步且时长不稳定的:用户可以接受排队和等待,但需要尽快知道“请求已受理、正在等待还是正在运行”,并持续看到后续进度或流式输出。

2. 候选方案

  • 直接同步调用:链路直观,但长任务、故障恢复和扩缩容都被耦合在请求路径中。
  • 进程内Event Bus与Redis Pub/Sub:适合即时通知和唤醒;订阅者离线时不承担可靠交付。
  • 数据库队列与轮询:事实一致性强,但领取、轮询、积压、顺序与并发控制需要自己实现。
  • 持久化消息系统:RabbitMQ、Kafka与RocketMQ。

3. RabbitMQ、Kafka与RocketMQ的取舍

  • RabbitMQ:复杂路由、逐条确认、优先级是强项;若要获得“同一对话有序、可回放”的语义,需要补充一致性哈希路由和单活消费等约束。
  • Kafka:按对话ID分区,分区内有序;消费者可以在自己的位点继续消费,天然容纳积压与重放。
  • RocketMQ:作为同类替代实现,内建延迟消息、事务消息、重试和死信等业务消息能力;本文不展开为实际落地方案。

4. Aviary为什么选Kafka

  • 把每段对话视为一条持续推进、可回放的事件流。
  • 同一对话串行,不同对话依靠分区并行。
  • 消息积压不是丢失,而是等待处理的工作;暂停或恢复后可以从消费位点继续。
  • 高峰时先由Kafka吸收突发请求,再由Worker按可用能力处理;用户侧不必等待任务完成才得到反馈,而是先展示已受理、排队或运行中的状态,再通过交付链路持续展示结果。
  • 事件流负责推进,数据库负责事实与查询,Redis负责快速唤醒和边缘扇出。

5. 可靠性不是Kafka自动给出的

  • 早期方案:业务事务提交后直接异步投递Kafka,留下“事实已保存、事件未发出”的窗口。
  • 演进方案:业务事实与Outbox在同一事务内冻结;发布器可靠投递,成功后再标记完成。
  • 消费端先落库、再提交位点;以event_id去重,接受at-least-once并把重复处理收敛为幂等。

6. 仍然需要保留的边界

  • 大对象、流式文本和随机查询不把Kafka当数据库使用。
  • 取消等抢占信号不应被长任务队头阻塞,需要独立的控制面路径。
  • Kafka只是内部工作主干,不承担用户侧实时展示的全部职责。