dodo-agent

流式控制与任务管理

在上节课中,我们已经实现了一个具备会话记忆、工具调用、联网搜索、参考来源以及推荐问题能力的智能问答系统。从功能层面来看,它已经能够完成完整的智能体对话流程。但如果将这样的系统直接上线,很快就会遇到一系列典型的工程问题。 例如:用户在…

TL;DR

在上节课中,我们已经实现了一个具备会话记忆、工具调用、联网搜索、参考来源以及推荐问题能力的智能问答系统。从功能层面来看,它已经能够完成完整的智能体对话流程。但如果将这样的系统直接上线,很快就会遇到一系列典型的工程问题。 例如:用户在…

在上节课中,我们已经实现了一个具备会话记忆、工具调用、联网搜索、参考来源以及推荐问题能力的智能问答系统。从功能层面来看,它已经能够完成完整的智能体对话流程。但如果将这样的系统直接上线,很快就会遇到一系列典型的工程问题。 例如:用户在回答生成过程中点击“停止生成”,系统是否真的停止了模型推理?如果用户连续点击发送消息,多个请求是否会同时执行?当用户关闭页面或网络断开时,后端任务是否仍在运行?如果大量长任务持续堆积,是否会拖垮整个系统?这些问题的本质并不在于大模型能力,而在于智能体系统的工程化能力。 一个真正可用的智能对话系统,必须能够对每一次执行任务进行完整的生命周期管理,包括任务注册、并发控制、流式终止以及资源释放。因此,在系统架构中我们引入了一个核心组件: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:流式输出发射器 - Disposable:模型流式控制句柄 - agentType:当前任务类型 源码中的定义如下:

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 则负责数据输出通道管理。二者结合,形成了一套完整的流式任务控制体系。 有了这一层能力之后,我们的智能问答系统不仅能够完成复杂的推理与工具调用,还能够在高并发环境下稳定运行,这也是一个生产级智能体系统必须具备的核心工程能力。

版本提示

模型、框架与接口会持续变化。涉及版本号、参数与生产配置时,请在实践前对照对应官方文档。

LLMentor系统化学习大模型应用工程

内容来自个人课程知识库备份,并经过结构化整理。技术版本持续演进,生产使用前请结合官方文档验证。