MQ 在 AI 后端的新角色:从异步解耦到 Agent 事件总线
MQ 的传统角色
MQ 在后端架构里一直是”老三样”:
- 异步解耦:订单创建后,发消息通知库存、物流、通知系统,各系统独立消费。
- 削峰填谷:秒杀流量来了,MQ 先扛着,后端慢慢消费。
- 最终一致性:分布式事务的兜底方案——本地事务 + MQ 消息 = 可靠异步通知。
这三个角色在今天仍然重要。但 AI 后端的出现,给了 MQ 一些新的可能性。
Agent 需要一个事件总线
看一个场景:用户在 Spring AI 应用里发起一个对话,Agent 需要串联多个操作——搜索 ES、查 MySQL、调用外部 API、生成回答。
传统做法是把这些调用写在 Agent 的逻辑里,串行执行。但如果 Agent 需要等待一个长时间操作(比如跑数据分析、等第三方回调),串行就不行了。
这时候 MQ 可以扮演 Agent 事件总线:
1 | 用户提问 → Agent 发布"需要搜索文档"事件 → MQ |
这不是传统 MQ 的”异步解耦”,而是一种”事件驱动的 Agent 工作流”。
Spring AI Alibaba Graph 内置了类似机制
Spring AI Alibaba 的 Graph 引擎本质就是一个事件驱动的 DAG(有向无环图)执行器。每个节点是一个独立步骤,节点之间通过事件传递来控制流转。
1 | var graph = new StateGraph("agent_workflow") |
这里的每个节点可以理解为消费某个”事件”并产生下一个”事件”。而 Graph 引擎负责调度和状态管理。如果你需要一个跨服务的 Agent 工作流,就可以把 Graph 节点替换为 MQ 消息。
MQ 在 AI 后端的具体新用法
1. Embedding 生成的异步流水线
文档写入后,需要调用 Embedding 模型生成向量。但 Embedding API 有 QPS 限制,同步调用会卡住业务接口。
1 | // 业务接口只发消息 |
2. 长时间 Agent 任务的异步回调
Agent 调用一个可能跑几分钟的分析任务,不能一直阻塞 HTTP 连接。
1 | Agent → MQ("run-analysis", {taskId, params}) |
这就需要 MQ 支持请求-响应模式或回调模式。RocketMQ 的 reply message 机制可以做到。
3. 多 Agent 间的消息传递
当多个 Agent 协作时(比如一个搜索 Agent + 一个分析 Agent + 一个回答 Agent),MQ 可以作为它们之间的通信基础设施。
这本质上是 Agent 间的发布-订阅模式。一个 Agent 完成自己的工作后,发布事件,其他关心这个事件的 Agent 自行消费。
MQ 选型考虑
| 特性 | 为什么对 Agent 重要 |
|---|---|
| 事务消息 | Agent 需要”做了一件事 + 通知别人”的原子性 |
| 顺序消息 | Agent 的步骤有顺序依赖 |
| 延迟消息 | Agent 需要在特定时间后执行某个操作 |
| 死信队列 | Agent 的任务失败了,需要兜底处理 |
RocketMQ 在这几个特性上比 RabbitMQ 和 Kafka 更适合 Agent 场景。RocketMQ 原生支持事务消息和延迟消息,而 RabbitMQ 需要插件、Kafka 需要外部协调。
意味着什么
- MQ 在 AI 后端的角色正在从”基础设施”升级为”架构中枢”。 不只是传输消息,而是编排 Agent 的工作流。
- 事件驱动架构和 Agent 架构天然契合。 Agent 的本质是”感知 → 决策 → 执行 → 反馈”的循环,而事件驱动是描述这个循环最自然的方式。
- RocketMQ 在 AI 后端里的价值可能被低估。 它的事务消息和延迟消息能力,在 Agent 工作流场景下比 Kafka 的纯高吞吐更有用。
一句话:MQ 不只是用来解耦订单和库存的。在 AI 后端里,它是 Agent 的事件总线——编排多步骤推理、管理长时间任务、连接多个 Agent 协作的中央管道。
