在上节课中,我们已经实现了一个具备会话记忆、工具调用、联网搜索、参考来源以及推荐问题能力的智能问答系统。从功能层面来看,它已经能够完成完整的智能体对话流程。但如果将这样的系统直接上线,很快就会遇到一系列典型的工程问题。 例如:用户在回答生成过程中点击“停止生成”,系统是否真的停止了模型推理?如果用户连续点击发送消息,多个请求是否会同时执行?当用户关闭页面或网络断开时,后端任务是否仍在运行?如果大量长任务持续堆积,是否会拖垮整个系统?这些问题的本质并不在于大模型能力,而在于智能体系统的工程化能力。 一个真正可用的智能对话系统,必须能够对每一次执行任务进行完整的生命周期管理,包括任务注册、并发控制、流式终止以及资源释放。因此,在系统架构中我们引入了一个核心组件:AgentTaskManager,用于统一管理所有智能体执行任务。
为什么必须管理流式任务?
在很多初期 Demo 中,开发者往往只关注一件事情:如何把模型生成的 token 实时推送到前端。例如通过 SSE 流,将模型输出逐 token 发送到页面。但在真实的工程系统中,仅仅做到“流式输出”远远不够,更重要的是如何管理这个流式任务本身的生命周期。 一个非常常见的问题是“停止生成”。很多早期系统的实现方式只是简单关闭前端的 SSE 连接:
前端点击停止生成
↓
服务端停止发送数据sinks.tryComplete
↓
大模型仍然在输出
表面上看,前端确实停止接收内容了,但模型推理其实仍然在后台继续执行。换句话说,流关闭了,但任务并没有停止。模型仍然在持续生成 token,只是这些 token 已经没有地方发送。 这种实现方式在真实生产环境中会带来两个严重问题: 第一是资源浪费。 大模型推理是典型的计算密集型任务,在 ToB 场景中,大多数企业使用的是私有化部署的大模型。一旦推理开始,就会持续占用 GPU 算力。如果任务没有被真正终止,这些算力仍然会被持续消耗,最终导致模型并发能力下降、token 输出速度变慢,用户体验明显变差。 第二是并发堆积。 如果系统中存在大量已经“无人接收”的后台任务,它们仍然会占用线程、连接以及推理资源。这些任务在系统中不断累积,就会形成典型的“幽灵任务”,最终导致系统吞吐能力下降,严重时甚至可能拖垮整个服务。 因此,在工程实践中,仅仅关闭前端连接是远远不够的。一个成熟的智能对话系统必须能够真正终止任务执行,并释放底层推理资源。这正是任务管理机制存在的意义。
Reactor Disposable流式控制
在我们的系统中,智能体流式输出是基于 Reactor 响应式编程模型实现的。当我们调用模型进行流式生成时,本质上是在创建一个 Reactive Stream。这个流在启动后会持续产生 token,并通过回调函数推送到下游。 例如:
Disposable
sink
一旦 subscribe() 被调用,整个流式任务就已经开始执行。模型会不断生成 token,并通过回调函数推送到 sink 中,再由 sink 转发到前端 SSE 通道。这里返回的 Disposable 对象非常关键。它并不是普通的数据对象,而是一个订阅关系的控制句柄。可以把它理解为当前流式任务的“遥控器”。只要这个对象存在,流就会持续运行;一旦调用 dispose(),Reactor 就会立即取消订阅关系,并向上游发送 cancel 信号,从而终止模型生成。
disposable
这一步的意义在于: 真正从源头终止模型推理,释放大模型资源,而不仅仅是停止向前端数据发送。
任务管理器
为了统一管理所有执行任务,我们在系统中实现了一个 AgentTaskManager。 其核心结构非常简单:使用一个线程安全的 ConcurrentHashMap 来记录当前系统中所有正在执行的任务。
private
每个任务对应一个 TaskInfo 对象,其中包含三个关键元素:
- Sinks.Many
public
}
这样,一个任务在系统中的结构就变得非常清晰:
conversationId
│
▼
├─ sink → SSE输出通道
├─ disposable → 模型执行控制
├─ agentType →
└─ createTime
通过这种结构,我们就可以同时控制: - 数据输出通道 - 模型推理任务 - 任务生命周期状态
任务注册与并发控制
在智能对话系统中,同一个会话在同一时间通常只允许存在一个执行任务。否则,如果用户在回答尚未完成时再次发送消息,就可能出现多个任务同时输出,导致流式输出的内容混乱。 因此在任务执行前,需要先检查是否已有任务在运行。 在 BaseAgent 中就实现了这样的检查逻辑:
protected
}
只有在确认当前会话没有运行任务时,系统才会注册新任务:
AgentTaskManager
任务注册逻辑在 AgentTaskManager 中实现:
public
log
taskMap
}
这样可以确保: - 同一会话不会并发执行多个任务 - 输出内容不会互相混乱 - 模型推理资源不会被重复占用
停止生成
当用户点击“停止生成”按钮时,系统会调用:
taskManager
在 AgentTaskManager 中,停止任务包含三个步骤。
终止模型推理
首先获取任务对应的 Disposable:
Disposable
if
disposable
}
这一步会触发 Reactor 的 cancel 信号,从而真正终止模型生成,从而释放大模型的算力资源。
关闭流式输出
接下来系统会向前端发送一条停止消息,并关闭流通道:
sink
sink
这样前端可以正确结束当前输出。
清理任务状态
最后,将任务从任务表中移除:
taskMap
这一步确保当前会话可以继续发起新的请求。
任务生命周期
在本系统中,智能体任务经历完整的生命周期:
用户请求
│
▼
任务注册
│
▼
模型推理执行React
│
▼
流式输出
│
├─ 正常完成
└─ 用户停止
│
▼
任务清理
在代码实现中,BaseAgent 会在任务结束时主动调用:
removeTask
从而保证任何情况下任务都不会残留在系统中。
总结
在智能对话系统中,流式输出只是表面能力,真正决定系统稳定性的,是任务生命周期管理机制。如果没有任务管理,停止生成只会关闭前端连接,模型仍然会继续运行,最终导致资源浪费和并发堆积。 通过引入 AgentTaskManager,系统能够对每一次智能体执行任务进行统一管理,实现任务注册、并发控制、流式终止以及资源释放。Disposable 作为 Reactor 流的控制句柄,使系统可以在任意时刻终止模型推理任务;而 Sinks.Many 则负责数据输出通道管理。二者结合,形成了一套完整的流式任务控制体系。 有了这一层能力之后,我们的智能问答系统不仅能够完成复杂的推理与工具调用,还能够在高并发环境下稳定运行,这也是一个生产级智能体系统必须具备的核心工程能力。