mistral.rs Rust SDK 并发请求批处理实战:多路并行请求与吞吐统计

mistral.rs Rust SDK 并发请求批处理实战:多路并行请求与吞吐统计 mistral.rs Rust SDK 并发请求批处理实战多路并行请求与吞吐统计【免费下载链接】mistral.rsFast, flexible LLM inference项目地址: https://gitcode.com/GitHub_Trending/mi/mistral.rs本文基于 mistral.rs 官方示例文档 batching.md 及其源码 mistralrs/examples/advanced/batching/main.rs讲解如何用 Rust SDK 以非流式方式同时发送多个推理请求让引擎自动将它们合并为批处理执行并逐请求观测 prefill / decode 阶段的 token 吞吐。读完后你将能够搭建一个带 ISQ 量化与分页注意力的模型实例、用join_all并发分发多个聊天请求、解析响应中的Usage统计结构并理解引擎侧的调度原理。为什么需要并发请求批处理在真实服务场景中多个用户请求几乎总是同时到达推理引擎。如果逐条串行发送GPU 大量时间花在等待单条序列的计算上硬件利用率很低。mistral.rs 的引擎内部带有调度器scheduler会把同一时刻在队列中的多个请求合并为同一个 batch 执行前向计算从而显著提升整体吞吐。这个示例要演示的正是应用层只需要把多个请求的 future 同时挂起批处理由引擎自动完成不需要手工组织 batch 张量。完整示例代码与执行方式文档给出的运行命令是cargo run --release --example batching -p mistralrs完整源码如下与 mistralrs/examples/advanced/batching/main.rs 一致//! Concurrent request batching by sending multiple requests in parallel. //! //! Run with: cargo run --release --example batching -p mistralrs use anyhow::Result; use mistralrs::{ ChatCompletionResponse, IsqBits, ModelBuilder, PagedAttentionMetaBuilder, TextMessageRole, TextMessages, Usage, }; const N_REQUESTS: usize 10; #[tokio::main] async fn main() - Result() { let model ModelBuilder::new(Qwen/Qwen3-4B) .with_auto_isq(IsqBits::Eight) .with_logging() .with_paged_attn(PagedAttentionMetaBuilder::default().build()?) .build() .await?; let messages TextMessages::new() .add_message( TextMessageRole::System, You are an AI agent with a specialty in programming., ) .add_message( TextMessageRole::User, Hello! How are you? Please write generic binary search function in Rust., ); let mut handles Vec::new(); for _ in 0..N_REQUESTS { handles.push(model.send_chat_request(messages.clone())); } let responses futures::future::join_all(handles) .await .into_iter() .collect::std::result::ResultVec_, _()?; let mut max_prompt f32::MIN; let mut max_completion f32::MIN; for response in responses { let ChatCompletionResponse { usage: Usage { avg_compl_tok_per_sec, avg_prompt_tok_per_sec, .. }, .. } response; dbg!(avg_compl_tok_per_sec, avg_prompt_tok_per_sec); if avg_compl_tok_per_sec max_prompt { max_prompt avg_prompt_tok_per_sec; } if avg_compl_tok_per_sec max_completion { max_completion avg_compl_tok_per_sec; } } println!(Individual sequence stats: {max_prompt} max PP T/s, {max_completion} max TG T/s); Ok(()) }模型构建四个关键 builder 方法ModelBuilder::new(Qwen/Qwen3-4B)自动检测架构的通用 builder模型权重从 Hugging Face 仓库下载.with_auto_isq(IsqBits::Eight)启用 in-situ 量化ISQ8-bit 量化在推理加载时把权重压到 8 位显著降低显存占用。自动模式会根据平台选择具体量化类型也可改用with_isq指定Q4_0、Q8_0、HQQ8等具体类型见 mistralrs/src/lib.rs 的 crate 文档.with_logging()开启日志输出便于观察调度过程.with_paged_attn(PagedAttentionMetaBuilder::default().build()?)启用分页注意力paged attention把 KV cache 组织为固定大小的页是支撑并发序列数较高的关键配置。这里示例发送N_REQUESTS 10个并发请求低于模型实例的默认调度容量——从源码看ModelBuilder的max_num_seqs默认值为 32见 mistralrs/src/auto_model.rs 中的字段默认值即引擎调度器默认按最多 32 条并发序列规划 KV cache 资源10 路并发可以完整落入同一批窗口内。并发分发future 挂起即入队核心并发逻辑非常简洁let mut handles Vec::new(); for for _ in 0..N_REQUESTS { handles.push(model.send_chat_request(messages.clone())); } let responses futures::future::join_all(handles).await...for循环中model.send_chat_request(messages.clone())返回的是尚未 await 的 future循环只是把这些 future 挂起并不会真正等待任何一个请求完成。随后futures::future::join_all(handles)同时驱动全部 10 个 future直到全部返回ChatCompletionResponse。这正是把请求塞给引擎、让引擎自己组 batch的应用层写法每条请求独立持有应答通道互不阻塞。从源码看send_chat_request的实现在 mistralrs/src/model.rs方法内部创建一个容量为 1 的 oneshot channel把消息、采样参数等包装成Request::Normal并通过self.runner.get_sender(model_id)?.send(request)投递给引擎然后循环接收结果跳过 agentic 进度类中间事件直到ResponseOk::Done携带最终响应返回。由于多个 future 在join_all中被并发 poll这些请求会几乎同时进入引擎侧队列随后被调度器合并成批执行。引擎侧多个请求如何变成一个 batch请求进入引擎后调度方式取决于 builder 是否启用了分页注意力。从 mistralrs/src/model_builder_trait.rs 的scheduler_config_from_pipeline可以看到两条路径启用 paged attn 时构造SchedulerConfig::PagedAttentionMeta其中max_num_seqs来自 builder默认 32并附带max_num_batched_tokens、max_prefill_chunk_tokens、max_decode_steps_before_prefill等默认常量用于控制单个迭代中允许处理的 token 上限与 prefill 分块未启用时退化为SchedulerConfig::DefaultScheduler { method: Fixed(max_num_seqs) }即固定上限的并发调度。引擎主循环mistralrs-core/src/engine/mod.rs每个调度周期从队列取出待处理序列区分 prompt 阶段序列与 decode 阶段序列把它们拼装成统一的 batch 交给 pipeline 前向计算。也就是说示例中 10 条序列会在同一个或少数几个调度迭代内共享一次 GPU 前向这就是批处理带来的吞吐增益来源。KV cache 的每序列耗时则分别累计在各自的 sequence 状态中total_prompt_time/total_completion_time见 mistralrs-core/src/sequence.rs最终换算为每请求的吞吐统计。解析 Usage每请求的吞吐指标示例解构响应时用到的Usage结构定义在 mistralrs-core/src/response.rs是 OpenAI 兼容 usage 的超集字段含义prompt_tokens/completion_tokens/total_tokens各阶段 token 计数prompt_tokens_details.cached_tokens命中前缀缓存、免于重算的 prompt token 数avg_tok_per_sec整个请求的平均 token 吞吐avg_prompt_tok_per_secprefill 阶段吞吐PP T/savg_compl_tok_per_secdecode 阶段吞吐TG T/stotal_time_sec/total_prompt_time_sec/total_completion_time_sec对应阶段的累计耗时秒吞吐换算在序列完成时计算以total_prompt_toks / total_prompt_time × 1000的形式把毫秒计时转为 token/s见 mistralrs-core/src/sequence.rs。示例最后打印的Individual sequence stats: {max_prompt} max PP T/s, {max_completion} max TG T/s即各请求中的最高 prefill 与最高 decode 吞吐。一个值得注意的细节示例代码统计max_prompt的条件写成了if avg_compl_tok_per_sec max_prompt { max_prompt avg_prompt_tok_per_sec; }——比较用的是 completion 速度、赋值用的却是 prompt 速度见 main.rs 第 55-61 行。按当前写法max_prompt实际输出的是在 completion 速度超过其当前值的那些请求中对应的 prompt 速度。如果你需要严格的最大 prefill 速度应把比较条件改为avg_prompt_tok_per_sec max_prompt。这是示例代码中的一个可改进点阅读输出结果时须知。可调参数与扩展方向N_REQUESTS并发请求数。调大可压出更高整体吞吐但受max_num_seqs默认 32与显存中 KV cache 容量约束序列过多时新请求会排队等待调度器腾出空位单请求延迟上升。with_paged_attn分页注意力让 KV cache 以页为单位分配是并发场景下显存效率与调度灵活性的基础建议保留。IsqBitsEight与Four直接决定权重显存占用与量化精度量化后 KV cache 可用空间也相应变化可结合显存余量选择。流式变体把send_chat_request换成model.stream_chat_request(messages.clone())即可让每个并发请求以流式返回增量 chunk返回的Stream实现futures::Stream并发 流式的组合方式不变。同步 API非 async 场景可用mistralrs::blocking::BlockingModel的send_chat_request其内部同样是block_on包装同一个请求路径见 mistralrs/src/blocking.rs。相关示例仓库中还有 多模型并发MultiModelBuilder单进程内多个模型共享一个引擎、分页注意力专项示例 与 流式生成可与本篇配合阅读。小结这个 batching 示例展示了 mistral.rs Rust SDK 的并发使用范式应用层用join_all并发挂起多个send_chat_requestfuture引擎内部的调度器自动把并发请求合并成批执行再由每请求的Usage结构给出 prefill/decode 吞吐。理解max_num_seqs默认容量32、paged attention 的角色以及Usage各字段的计算方式就能在此基础上正确评估并发配置的收益并规避示例代码中统计条件的一处笔误。【免费下载链接】mistral.rsFast, flexible LLM inference项目地址: https://gitcode.com/GitHub_Trending/mi/mistral.rs创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考