dodo-agent

流式响应的一些重要概念

在前面的学习中,我们已经实现了一个具备 流式输出、任务管理、停止生成 等能力的完整的智能对话系统。在这个过程中,你会发现代码里频繁出现一些概念,例如: Mono Flux Sinks Disposable flatMap 背压(on…

TL;DR

在前面的学习中,我们已经实现了一个具备 流式输出、任务管理、停止生成 等能力的完整的智能对话系统。在这个过程中,你会发现代码里频繁出现一些概念,例如: Mono Flux Sinks Disposable flatMap 背压(on…

在前面的学习中,我们已经实现了一个具备 流式输出、任务管理、停止生成 等能力的完整的智能对话系统。在这个过程中,你会发现代码里频繁出现一些概念,例如: - Mono - Flux - Sinks - Disposable - flatMap - 背压(onBackpressureBuffer) 对于很多第一次接触的同学来说,这些概念一开始确实会有些陌生,甚至会感觉非常复杂,难以理解。但是在大模型应用开发中,响应式肯定是绕不过去的,这是提升用户体验最重要的手段,实际上,大部分情况下,只要搞懂下面几个核心概念,就已经足够理解绝大多数流式代码的运行方式: - Mono:表示单值异步结果 - Flux:表示连续的数据流 - Sinks:手动向流中发送数据 - Disposable:控制流式任务的执行与停止 - flatMap:在流中执行异步任务 - Backpressure(背压):控制数据流速度 接下来,我们结合实际代码逐个理解这些概念。

Mono:单值异步结果

在传统的 Spring MVC 项目中,一个方法通常会直接返回结果:

public


}

这里的 String 就是最终结果。 但在响应式编程中,方法不会直接返回结果,而是返回一个未来才会产生结果的对象:

public

}

Mono 可以理解为: 未来某个时间点会返回一个 String 也就是说: | 类型 | 含义 | | --- | --- | | String | 已经得到结果 | | Mono | 未来会得到结果 |

在大模型开发中,如果是响应式项目,并且调用的是 非流式接口(一次性返回完整回答),通常就会使用 Mono。 例如:

Mono<String> result = llmClient.call(prompt);

Mono:启动异步任务

除了表示“单值结果”,Mono 还经常被用来 启动异步任务,这一点在智能体开发中非常常见。 例如:

Mono

}

这里的含义是: - fromRunnable 将一个普通任务包装成 Mono - subscribe() 时任务开始执行 - 任务在后台线程运行 如果再配合线程调度器:

Mono

}
.
.

执行流程就是:

创建 Mono 任务
      ↓
subscribe 触发执行
      ↓
任务在线程池运行
      ↓
执行 Agent 逻辑

这种写法在 智能体系统中非常常见,因为很多任务并不是简单返回一个值,而是: - 执行工具调用 - 查询搜索引擎 - 推送流式结果 - 写入 Sink 例如在智能体实现中,很可能会出现类似逻辑:

Mono
    agent

}

这里 Mono 的作用其实是:用响应式方式启动一个后台任务,让智能体逻辑在异步线程中运行。

Flux:连续的数据流

如果 Mono 表示一个结果,那么 Flux 就表示一串连续产生的数据。 在大模型流式输出场景中,模型通常不会一次返回完整文本,而是逐步生成 token,例如:

Hel
lo
 wor
ld

这种持续产生数据的过程,就可以用 Flux 表示:

Flux

然后系统可以订阅这个流:

stream

}

在你的系统中,token 会被推送到前端:

Disposable

            sink

整个流程可以理解为:

大模型

因此: | 返回类型 | 使用场景 | | --- | --- | | Mono | 单次返回 | | Flux | 流式输出 |

Sinks:手动向流中发送数据

在很多情况下,数据流并不是直接来自某个 API,而是需要系统自己向流中写入数据。 例如在你的代码中:

Sinks

这里创建了一个 可以主动写入数据的流通道。 当模型生成 token 时:

sink

就可以把 token 推送到前端,最后再把 sink 转换成 Flux 返回:

return

整个数据链路就变成:

LLM token
   ↓
sink
   ↓
Flux
   ↓
SSE
   ↓
前端页面

可以这样理解: | 组件 | 作用 | | --- | --- | | Sink | 生产数据 | | Flux | 数据流 | | Subscriber | 消费数据 |

Disposable:流的控制器

当我们订阅一个 Flux 时,会得到一个 Disposable:

Disposable
        chatModel

                    sink

这个对象可以理解为:当前流式任务的控制开关。 如果不做任何操作,流会一直运行直到模型生成结束。但如果用户点击 停止生成,就可以调用:

disposable

这一步非常关键,它会: - 取消 Flux 订阅 - 停止 token 生成 - 终止后台推理任务 这也是你在 AgentTaskManager 中保存 Disposable 的原因:

taskManager

当用户停止生成时:

disposable

任务就会被真正终止。

flatMap:在流中执行异步任务

在响应式编程中,经常需要在数据流中执行新的异步操作,这时候就会用到 flatMap。 例如:

Flux

Flux
        llmClient
)

这里的执行逻辑是:

Flux
        ↓
flatMap
        ↓
调用大模型
        ↓
返回多个结果

flatMap 的核心特点是: 它可以把一个元素转换成一个新的异步流,并合并结果。 在智能体系统中,常见的场景包括: - 查询多个工具 - 并发调用多个 API - 多个子任务并行执行 例如:

问题
 ↓
拆分任务
 ↓
flatMap并发执行
 ↓
汇总结果

Backpressure:背压机制

在流式系统中,还有一个非常重要的概念:背压(Backpressure)。 背压解决的是这样一个问题: 如果数据生产速度 > 数据消费速度,会发生什么? 例如:

模型每秒生成 100 token
前端只能处理 20 token

如果不做控制: - 内存会不断增长 - 数据会积压 - 系统最终可能崩溃 因此响应式框架提供了背压机制。在代码中:

Sinks

onBackpressureBuffer() 的意思就是: 如果消费不过来,就先放到缓冲区。 常见背压策略包括: | 策略 | 含义 | | --- | --- | | Buffer | 缓存数据 | | Drop | 丢弃数据 | | Latest | 只保留最新数据 |

在大模型系统开发中,Buffer 是最常见的策略。

总结

如果从智能体系统的角度来看,其实只需要掌握下面这些概念,就已经足够理解绝大多数代码。 | 概念 | 含义 | 在大模型系统中的作用 | | --- | --- | --- | | Mono | 单值异步结果 | 非流式调用/启动异步任务 | | Flux | 多个数据流 | token 流式输出 | | Sinks | 手动发送数据 | 推送 token 到 SSE | | Disposable | 流控制句柄 | 停止生成 | | flatMap | 异步流转换 | 并发调用工具或模型 | | Backpressure | 背压控制 | 防止 token 堆积 |

在大模型应用开发中,响应式编程几乎已经成为一种非常常见的技术选择。它不仅可以支持 Token 级流式输出,还能够很好地处理 高并发任务、异步工具调用以及任务生命周期管理 等复杂场景。

版本提示

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

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

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