使用 Rust 开发 AI Agent(第二讲):结构化输出、流式调用与并发限流控制
本讲聚焦于 AI Agent 开发中三个关键的基础能力:通过 JSON Schema 约束大语言模型实现类型安全的结构化输出、使用流式(Streaming)响应改善用户体验,以及利用信号量与指数退避重试机制实现高效的并发调用与限流控制。
核心要点
- 大语言模型默认返回自由格式的自然语言,对程序解析极不友好;结构化输出(Structured Output)通过预先指定的 JSON Schema 约束模型输出格式,保证返回结果可直接反序列化使用。
- 在 Rust 中,结构化输出的完整链路为:使用
serde定义数据结构 → 通过schemars的JsonSchemaderive 宏在编译期自动生成 JSON Schema → 传给大语言模型 → 响应返回后直接反序列化。整个过程类型安全,格式不匹配在编译期即可发现。 - 结构化输出在工具调用(Tool Calling)中至关重要——大语言模型必须以精确的结构化数据告知运行时“调用哪个工具、传什么参数”,不能是自由格式的文字。
- 非流式请求需要等待模型一次性生成完整响应,等待时间长、用户体验差;流式响应通过
async_stream逐块接收模型生成的增量内容。 - 当任务数量庞大(如后续测试集有 100 多道题)时,顺序调用耗时不可接受(可能超过 1000 秒),必须并发执行;但并发过高会触发 API 提供商的速率限制(Rate Limit),因此需要通过信号量精确控制并发上限。
- 偶发的限流错误可采用指数退避重试策略处理:第一次失败等 1 秒重试,第二次等 2 秒,第三次等 4 秒(默认策略),能有效应对大部分偶发性 API 错误。
详细解析
一、结构化输出(Structured Output)
1. 要解决的问题:自然语言输出对程序不友好
大语言模型默认返回的是自然语言格式,对人类友好但程序无法直接处理。虽然可以尝试通过模式匹配、关键字提取等方式解析信息,但由于文本过于自由,总有遗漏的情况。每次解析都是一次痛苦的尝试。
结构化输出的做法是:预先指定一个 JSON 的模式,在调用大语言模型时将该模式传给模型。模型会严格按照该格式返回,不会有任何自由发挥。开发者收到的直接就是合法的 JSON,反序列化即可使用。
2. 为什么在工具调用中至关重要
大语言模型在工具调用场景下需要告诉运行时:
- 调用哪个工具
- 传什么参数
这些信息必须是精确的结构化数据。如果模型以自由文本返回这些指令,运行时将无法可靠地解析和执行。
3. Rust 中实现结构化输出的技术栈
| 组件 | 作用 |
|---|---|
serde | 定义 Rust 数据结构,提供序列化与反序列化能力 |
schemars(JsonSchema derive) | 在编译期扫描结构体的字段和类型信息,自动生成 JSON Schema |
response_format(API 参数) | 将生成的 JSON Schema 传给大语言模型,约束其返回格式 |
代码结构示例:
- 新建
models模块下名为action_plus.rs的子模块,用于存放要返回的数据结构定义。 - 在
main.rs中声明模块(mod models;及对应的子模块声明)。 - 核心结构体上添加
#[derive(JsonSchema, Serialize, Deserialize)]等 derive 宏。 - 为演示目的,数据结构中包含另一个结构体和一个枚举,这些也需要添加对应的 derive 宏。
4. 在请求中配置结构化输出
在函数内部(请求发出之前),添加以下关键步骤:
- 生成 JSON Schema:使用
schemarscrate 中的宏schema_for!,传入自定义结构体ActionPlan。该宏在编译时扫描结构体的字段类型信息,生成一个schemars::schema::Schema类型的对象(注意:仍是 Rust 类型)。 - 转换为 JSON:将上述 Schema 对象转化为 JSON 格式。
- 构建 `response_format`:通过
response_format相关类型(create_response_format或等价方式),构造包含四个字段的配置:
description:对结构化结果的描述name:名称为action_planschema:传入生成的 JSON Schemastrict:设为true
- 严格模式(Strict Mode)说明:当
strict设为true时,所有字段都是必填(required)的,且不允许模型输出 schema 中没有的额外字段。如果不开启严格模式,大语言模型可能会自行添加一些奇怪的额外字段。 - 传入请求:在
request的messages下方配置response_format字段,传入构建好的 format 设置。 - 处理响应:得到响应后,从返回中取出字符串内容,用
serde_json反序列化为ActionPlan类型,返回给调用方(同时修改函数的返回类型)。
5. 执行效果
运行后,返回结果分为两层:
- 原始数据结构(模型未解析前的响应)
- 提取出的
ActionPlan对象,其中包含目标描述、步骤列表(示例中分为 9 步)、难度等级(简单)和预计时长(30 分钟)
二、流式调用(Streaming)
1. 为什么需要流式调用:非流式请求的等待问题
非流式调用(如上一节使用的普通 complete 方式)的流程是:
- 向大语言模型发出请求
- 等待模型一次性将完整响应返回
问题在于中间的等待过程有时较长,尤其是复杂问题。期间用户看不到任何中间反馈,对用户不够友好。
2. 流式调用的核心机制
流式调用通过 stream 方式获得模型响应,模型逐块、逐段地生成并返回内容,用户可以实时看到输出过程,体验更接近打字机效果。
技术栈
async_streamcrate(版本 0.3.6):提供stream!宏,用于编写异步流
核心代码逻辑
关键修改点在函数中请求发送之后的处理部分:
- 创建流式请求:使用
client.stream(...)(等价写法为client.create_stream(...))替代普通请求,然后.await。返回的类型是一个流,而不是一次性完整的响应。 - 处理流的语法:使用
while let循环 +stream.next().await持续等待下一个数据块到达。 - 逐块处理:每次流产生一个数据块,用
match判断是否为Ok:
- 如果
Ok,处理逻辑与之前非流式类似——从结果中取出delta。delta即“增量内容”,每次只包含新生成的那一小段文字,不包括之前已输出的内容。从delta中提取content字段,将增量文本yield Ok(...)返回到输出流。 - 如果不是
Ok,则yield Err(...)。这通常是由网络抖动或 API 端问题导致的。
- 流结束:当模型生成完毕、流耗尽时,
while let循环自动退出;随后stream!宏代码块结束,整个流结束。
函数签名及返回类型
函数返回类型变为:
impl Stream<Item = Result<...>>,要求实现Streamtrait- 由于还涉及异步操作,需加上
Future约束(impl Future<Output = ...> + Send等组合写法) - 最外层使用
stream!宏包裹异步代码块
注意事项
若返回的 Stream 类型未实现 Send,在异步运行时中可能报错;需确认引入的 Stream 类型来自 async_stream 并满足 Send 约束。
3. 客户端消费流:实测示例
在项目根目录创建 examples 文件夹(Rust 约定目录),将 s2c 目录下的某个 .rs 文件复制过去并改名为 stream_try.rs。同时需要在 Cargo.toml 中声明该示例:
[[example]]
name = "stream_try"
path = "examples/stream_try.rs"具体实现要点
- 固定内存地址:由于流返回过程中可能包含指向自身的内部指针,必须使用
Box::pin将它固定(pin)在栈内存上(示例中使用tokio::pin!宏)。 - 准备输出缓冲:创建一个可变的字符串
output(let mut output = String::new())。 - 消费流:用
while let循环不断取出流的元素。每个元素是一个Result:
- 若不是
Ok,返回错误(return Err(...)) - 若是
Ok,把其中的增量字符串拼接到output中,同时打印出来,形成动态输出效果
- 最终输出:等流耗尽后,打印拼好的完整结果。
调用方代码效果
向模型提问“道德经的第四章是什么内容?”,运行时可以看到:
- 模型生成过程中文本逐字逐句地动态输出
- 最终完整内容拼接成功后再次打印全文
体验显著优于非流式请求,因此本系列视频教学的重心确定为使用 stream 方式调用。
三、并发调用与限流控制
1. 何时需要并发
当需要同时处理多个任务时,例如:
- 后续将运行测试集,第一个级别就有 100 多道题
- 多个 Agent 并行工作
- 多个提示词同时请求
如果顺序调用:
- 每个请求快的约 1 秒,慢的可能超过 10 秒
- 100 多道题可能总计需要 1000 多秒,等不起
理论上并发执行可以把总时间压缩到约 10 秒。
2. 为什么需要限流控制
并发太高会带来新问题:AI API 提供商通常有速率限制(Rate Limit)。当同时发出过多请求时,超出的请求会被限流,返回错误。因此,必须精确控制并发的上限。
方案概述:
| 机制 | 作用 |
|---|---|
| 信号量 | 设置最大并发数(如 5 或 10)。即使创建了 100 个任务,也只有 5 或 10 个在同时运行,其余排队等候 |
| 指数退避重试 | 解决偶发的限流错误:第一次失败等 1 秒重试,第二次等 2 秒,第三次等 4 秒 |
在 Rust 生态中,这些功能可以较方便地实现。
3. 实现指数退避重试
添加 crate:backon。
为什么需要包一层闭包
重试逻辑需要能反复调用同一个操作,因此传入的参数不是 Future 本身,而是返回 `Future` 的闭包(`FnMut`)。原因是:
Future只能await一次,await完成后即被消耗- 闭包每次调用都会生成一个新的
Future,因此重试机制才能在失败之后重新执行
代码结构要点
async fn try_stream_with_retry(...) -> Result<String, ...> {
let op = || async { // 闭包内部封装实际请求逻辑
// 与之前流式请求逻辑相同的代码
// 最终返回 Ok(output) 或 Err(...)
};
// 执行重试
op.retry(...).await
}- 闭包内部代码与前面实现的流式请求逻辑一致
- 最后返回
Ok(output)(output为拼接好的完整字符串)
默认重试策略
backon 采用指数级默认策略:
- 第一次失败:等待 1 秒
- 第二次失败:等待 2 秒
- 第三次失败:等待 4 秒
- 最多重试 3 次(加第一次执行,最多运行 4 次)
- 如果 3 次重试全部失败,将最后一次的错误返回给调用方
4. 实现全局限流(信号量)
单独新建一个 .rs 文件(如 limit.rs)并在模块中注册。由于是教学演示,统一使用一个简单的全局限流设定(实际生产环境应根据不同 API 提供商的限制分别配置)。
使用的类型
| 组件 | 来源 | 作用 |
|---|---|---|
OnceLock | Rust 标准库 | 只能初始化一次的容器,用于存储全局单例 |
Semaphore | tokio::sync | 异步信号量,控制同时最多可执行的任务数 |
实现逻辑
static SEMAPHORE: OnceLock<Semaphore> = OnceLock::new();
pub async fn acquire_permit() -> ... {
// 初始化:创建最多允许 3 个并发任务的信号量
SEMAPHORE.get_or_init(|| Semaphore::new(3));
// 获取许可
SEMAPHORE.get().unwrap().acquire_owned().await
}变量的生命周期为整个程序('static),因此使用 static 存储。初始化的信号量最多允许 3 个任务同时持有许可。
5. 综合示例代码解读
在示例文件(stream_try.rs)中演示完整流程。先准备多个提示词,然后:
- 使用 `JoinSet`(`tokio::task::JoinSet`):一个任务容器,可向其中扔多个任务并发执行。每次
join_next()会等一个任务完成并取出其结果。 - 遍历提示词:对每个提示词:
- 创建一个
tracing::span,打上标签(日志中可看到是哪个提示词) - 使用
join_set.spawn(...)创建任务,任务体为async move块
- 任务内部先获取信号量许可:调用上述限流函数请求许可。由于信号量上限为 3,同一时刻最多 3 个任务能持有许可。如果请求不到,任务在此挂起等待,直到其他任务释放许可。
- 执行实际调用:获取到许可后,调用带重试功能的函数(
try_stream_with_retry),传入提示词。 - 手动释放许可:结果出来后(代码中所示范的写法)手动将许可还回信号量(即
drop(permit))。注释指出,这一步实际上不写也行(释放可通过drop隐式完成,但明确释放更清晰)。 - 返回结果:将
(prompt, output)一起返回,外层join时可以拿到对应关系。 - 消费 JoinSet 结果:使用
while let Some(res) = join_set.join_next().await循环:
- 每次等一个任务完成并拿结果
- 结果为嵌套
Result(join_next的Result和任务内部的Result),分别处理 - 成功时用
tracing打印提示词和结果 - 失败时用
tracing::error!记录 - 当所有任务跑完、
JoinSet返回None时,循环退出
最终运行效果是通过 tracing 输出各任务结果。作者提示打印过程可能需要优化,建议用 tracing 替代 println!。
方法与步骤
结构化输出实现步骤
- 安装所需 crate:
schemars(对应schema_for!宏) - 在
lm目录下复制complete.rs,重命名为structure.rs - 新建子模块文件(如
models/action_plus.rs),定义目标数据结构:
- 核心结构体(如
ActionPlan) - 辅助结构体与枚举(演示用)
- 均添加
#[derive(JsonSchema, Serialize, Deserialize)]
- 在模块声明处注册新子模块
- 修改
structure.rs中的函数:
- 函数名加后缀(如
try_complete_structure) - 保留原有 system/user message 配置
- 请求前用
schema_for!生成 schema → 转为 JSON - 构建
response_format,设置description、name、schema、strict: true - 请求中传入
response_format - 响应后反序列化字符串为结构体,更新函数返回类型
- 修改
main.rs调用方并运行验证
流式调用实现步骤
- 添加 crate:
async_stream(0.3.6 版本) - 复制
complete.rs为stream.rs,注册模块 - 修改返回类型:
impl Future包裹impl Stream<Item = Result<...>>,确保Send - 添加
stream!宏包裹逻辑 - 将普通请求改为流式请求(
stream()/create_stream()) - 使用
while let+stream.next().await循环取块 - 处理每个块:
Ok时取delta.content,yield Ok(...)- 出错时
yield Err(...)
- 创建 examples 目录和示例文件,在
Cargo.toml中声明[[example]] - 在示例中
Box::pin流、拼接增量文本、打印过程、输出完整结果 - 运行
cargo run --example <名称>验证
并发限制与重试实现步骤
- 添加 crate:
backon - 添加具备重试能力的流式函数(如
try_stream_with_retry):
- 将原请求逻辑包进闭包(返回
Future) - 调用
.retry(...)施加指数退避策略
- 新建限流模块文件,使用
OnceLock+tokio::sync::Semaphore创建全局信号量(示例上限为 3) - 主调用处:
- 创建
JoinSet - 遍历提示词,创建带
tracingspan 的任务 - 任务内:获取信号量许可 → 执行带重试的调用 → 释放许可 → 返回
(prompt, result) - 用
while let循环消费JoinSet结果,分别处理成功失败
案例与数据
结构化输出示例
- 提问目标返回格式:一个
ActionPlan结构,包含目标(goal)、步骤列表、难度、预计时间等字段 - 实际输出:目标明确,共分 9 步,难度为“简单”,预计 30 分钟完成
流式调用示例
- 问题:“道德经的第四章是什么内容?”
- 运行过程:文本增量逐段打印,产生打字机式动态输出效果;结束后打印完整拼接文本
并发需求数据
- 测试集规模:第一个级别就有 100 多道题
- 顺序调用耗时预估:
- 快者约 1 秒/题
- 慢者可能超过 10 秒/题
- 总计可能超过 1000 秒
- 并发预期:理论上可压缩到约 10 秒
- 信号量上限设定:3(教学演示)、5 或 10(实际场景可调)
指数退避默认参数
| 重试次数 | 失败后等待时间 |
|---|---|
| 第 1 次失败后 | 1 秒 |
| 第 2 次失败后 | 2 秒 |
| 第 3 次失败后 | 4 秒 |
- 最大重试次数:3 次;加上首次执行,最多运行 4 次
- 全部失败时:返回最后一次错误
限制与待确认问题
- 国内模型的兼容性问题:结构化输出方面,国内一些模型(如 DeepSeek)可能不支持使用 JSON Schema 的方式,需要调整代码实现。演示者明确表示“他使用 JSON Schema 好像是不行,这个你再探索一下”——该问题需读者针对具体模型自行验证和适配。
- 限流参数的简化:本次教学实现中限流是统一设定的单一值(上限 3),未按不同 API 提供商分别精细化配置。实际生产环境应根据各 API 提供商的速率限制做差异化设置。
- 打印优化提醒:演示过程中并发结果打印可能不够理想,演示者建议使用
tracing替代println!改善输出。 - 手动释放许可:代码中手动归还信号量许可属于示范写法;实际上不显式编写释放逻辑也可以(许可会在离开作用域时被丢弃并自动释放),但演示代码保留了显式释放的写法。
应用要点
- 实际项目建议按 API 提供商分别设定限流配置,不能一刀切使用固定并发数。
- 流式调用是本系列后续教学的默认调用方式,它比非流式更适合真实产品中的交互体验。
- 遇到偶发的限流错误时优先使用指数退避重试,通常可以解决大部分 API 偶发错误。
- 面向不同模型做兼容处理:在结构化输出特性上,不同厂商模型能力不完全一致,应针对使用的具体模型做相应适配(如 DeepSeek 可能的 JSON Schema 支持问题)。
总结
本讲完整覆盖了消息驱动的 AI Agent 开发中的三个关键工程问题:
- 结构化输出:通过
serde+schemars定义类型安全的输出契约,配合 API 的response_format(含严格模式)让模型精确返回 JSON,避免脆弱的文本解析。这对工具调用可靠性的意义重大。 - 流式调用:使用
async_stream的stream!宏改造请求方式,通过while let逐块消费增量内容并转发给上层,改善延迟体验。客户端层面通过固定内存地址(pinning)和增量拼接完成消费。 - 并发与限流:以
JoinSet组织多任务并发执行;以OnceLock+ 信号量(Semaphore)精确控制同时运行的请求数量;以backon库实现指数退避重试应对偶发性的 API 限流错误。
三个能力将在后续的测试执行场景(如一百多道题的测试集)中组合运用,实现高效且稳定的批量请求调度。