Firestore数据建模最佳实践:从Awesome Firebase精选案例中学到的经验
mistral.rs Python SDK 流式输出实战:异步迭代、FastAPI 集成与流中断错误处理
本篇指南围绕 mistral.rs Python SDK(mistralrs)的流式(streaming)输出能力展开,覆盖三个核心场景:在异步代码中消费流式响应、把流式生成接入 FastAPI 等 Web 框架、以及流中途失败时的错误处理。读完本文,你将掌握 stream=True 背后的数据模型与底层实现(ChatCompletionStreamer 的迭代协议与错误映射),并能写出可复制、可运行的异步流式与 Web 流式代码,同时理解生产环境下改用 HTTP 服务的原因。
前置:流式输出的基础用法
流式(streaming)是指模型按 token 逐块返回增量结果,而不是等完整回复生成完毕后一次性返回。在 mistral.rs Python SDK 中,入口是 Runner.send_chat_completion_request,它接收一个 ChatCompletionRequest,当其中 stream=True 时返回一个同步迭代器,而不是单一的响应对象。
from mistralrs import Runner, Which, ChatCompletionRequest
runner = Runner(
which=Which.Plain(model_id="Qwen/Qwen3-4B"),
in_situ_quant="4",
)
stream = runner.send_chat_completion_request(
ChatCompletionRequest(
model="default",
messages=[{"role": "user", "content": "Write me a haiku about ownership."}],
max_tokens=128,
stream=True,
)
)
for chunk in stream:
delta = chunk.choices[0].delta.content
if delta:
print(delta, end="", flush=True)
print()
这段基础用法取自 getting started 指南,对应的完整可运行示例见 examples/python/streaming.py(该示例使用 GGUF 模型加载方式,并演示了 presence_penalty、top_p、temperature 等采样参数与 stream=True 的组合使用)。
在写任何进阶用法之前,需要先理解两个关键事实:
- SDK 不提供原生异步迭代器:
send_chat_completion_request返回的是一个实现了 Python 迭代协议的对象,for chunk in stream中的next()调用是阻塞的。想在async代码中使用,必须自行包装(下文详述)。 - 每个 chunk 都是 OpenAI 兼容的流式结构:
choices[0].delta.content携带一段增量文本,它可能是None——例如携带finish_reason的最后一个 chunk。因此所有示例都先判断if delta:再打印。
流式 chunk 的数据模型
从 mistralrs.pyi 的类型声明可以看到流式响应的完整结构:
ChatCompletionChunkResponse:顶层对象,字段包括id、created、model、system_fingerprint、choices,以及可选的usage、adapter_generation、session_id;ChunkChoice:每个 choice 包含finish_reason(str | None)、index、delta和logprobs;Delta:增量内容,content(str | None)、role、可选的tool_calls与reasoning_content。
@dataclass
class Delta:
content: str | None
role: str
tool_calls: list[ToolCallResponse] | None = None
reasoning_content: str | None = None
@dataclass
class ChunkChoice:
finish_reason: str | None
index: int
delta: Delta
logprobs: ResponseLogprob | None = None
@dataclass
class ChatCompletionChunkResponse:
id: str
choices: list[ChunkChoice]
created: int
model: str
system_fingerprint: str
object: str
usage: Usage | None = None
adapter_generation: str | None = None
session_id: str | None = None
finish_reason 为空(None)时表示生成仍在进行;当某个 chunk 的所有 choice 都带有 finish_reason 时,迭代器就会终止。max_tokens 截断时 finish_reason 为 "length",正常结束时为 "stop" 等取值——判断循环是否结束,应依据 finish_reason 而非 delta.content 是否为 None。
底层原理:ChatCompletionStreamer 如何工作
流式迭代器的实现位于 mistralrs-pyo3/src/stream.rs 的 ChatCompletionStreamer。它内部持有一个 Tokio 异步通道的接收端 rx: Receiver<Response> 和一个 is_done 标志,通过 __iter__ / __next__ 暴露给 Python。
#[pyclass]
pub struct ChatCompletionStreamer {
rx: Receiver<Response>,
is_done: bool,
}
#[pymethods]
impl ChatCompletionStreamer {
fn __iter__(this: PyRef<'_, Self>) -> PyRef<'_, Self> { this }
fn __next__(mut this: PyRefMut<'_, Self>) -> Option<PyResult<ChatCompletionChunkResponse>> {
if this.is_done {
return None;
}
loop {
let recv_result = py.allow_threads(|| rx.blocking_recv());
match recv_result {
Some(resp) => match resp {
Response::AgenticToolCallProgress { .. } => continue,
Response::AgenticToolApprovalRequired { .. } => continue,
Response::BlockDenoisingProgress(_) => continue,
Response::File(_) => continue,
Response::ModelError(msg, _) => {
return Some(Err(PyValueError::new_err(msg.to_string())));
}
Response::ValidationError(e) => {
return Some(Err(PyValueError::new_err(e.to_string())));
}
Response::InternalError(e) => {
return Some(Err(PyValueError::new_err(e.to_string())));
}
Response::Chunk(response) => {
if response.choices.iter().all(|x| x.finish_reason.is_some()) {
this.is_done = true;
}
return Some(Ok(response));
}
_ => unreachable!(),
},
None => {
return Some(Err(PyValueError::new_err(
"Received none in ChatCompletionStreamer".to_string(),
)));
}
}
}
}
}
这段实现回答了文档中几个关键行为"为什么是这样":
- 迭代器会阻塞等待:
__next__在allow_threads中调用blocking_recv(),意味着每次next()都会挂起当前线程直到引擎产生下一个响应。这正是异步代码中需要把它丢进 executor 的原因。 - 错误以
ValueError形式作为下一个迭代项返回:引擎侧的ModelError(如生成中途 OOM)、ValidationError、InternalError都被映射为PyValueError。注意它不是立即抛出,而是作为迭代的"下一个元素"返回,因此已产生的 chunk 不受影响,部分输出得以保留。 - 工具进度事件被透明跳过:
AgenticToolCallProgress、AgenticToolApprovalRequired、BlockDenoisingProgress、File等内部事件在循环里continue掉,Python 侧只会看到纯内容 chunk。
请求侧的分发逻辑在 mistralrs-pyo3/src/lib.rs 的 send_chat_completion_request 中:它创建一个容量为 10_000 的 channel,将请求发送给引擎,随后根据是否为流式请求返回 Either<ChatCompletionResponse, ChatCompletionStreamer>——非流式拿到完整响应,流式拿到 ChatCompletionStreamer::from_rx(rx)。流式与同步等待的分流实现在 util.rs 的 send_request_with_optional_stream:流式请求直接返回接收端,非流式请求则循环接收直到拿到最终响应。
异步流式:用 executor 包装同步迭代器
SDK 的流式迭代器是同步阻塞的,因此要在 asyncio 中使用,标准做法是把 next(stream) 调用提交给事件循环的默认执行器,让它在线程池中执行,从而不阻塞事件循环。官方推荐写法如下:
import asyncio
from mistralrs import Runner, Which, ChatCompletionRequest
runner = Runner(Which.Plain(model_id="Qwen/Qwen3-4B"))
async def stream_response(prompt: str):
stream = runner.send_chat_completion_request(
ChatCompletionRequest(
model="default",
messages=[{"role": "user", "content": prompt}],
stream=True,
)
)
loop = asyncio.get_event_loop()
while True:
chunk = await loop.run_in_executor(None, next, stream, None)
if chunk is None:
break
delta = chunk.choices[0].delta.content
if delta:
yield delta
消费端用一个异步生成器即可,配合 async for 使用:
async def main():
async for delta in stream_response("Write a haiku."):
print(delta, end="", flush=True)
要点拆解:
loop.run_in_executor(None, next, stream, None):None是传给next的默认值参数,当迭代器耗尽时返回None而不是抛StopIteration,从而优雅退出循环;- 每次
await都会释放事件循环控制权,让其他协程(例如同时处理多个用户的请求)得以调度; stream_response是异步生成器,async for消费时逐块yield增量文本,非常适合进一步对接 WebSocket、SSE 等传输层。
流式接入 Web 框架:FastAPI 实战
同样的模式可以直接作为 FastAPI 的响应生成器使用。核心思路不变:send_chat_completion_request 返回同步迭代器,把它包进一个普通生成器函数,再交给 StreamingResponse。
from fastapi import FastAPI
from fastapi.responses import StreamingResponse
from mistralrs import Runner, Which, ChatCompletionRequest
app = FastAPI()
runner = Runner(Which.Plain(model_id="Qwen/Qwen3-4B"))
@app.get("/stream")
async def stream(prompt: str):
def iter():
s = runner.send_chat_completion_request(
ChatCompletionRequest(
model="default",
messages=[{"role": "user", "content": prompt}],
stream=True,
)
)
for chunk in s:
delta = chunk.choices[0].delta.content
if delta:
yield delta
return StreamingResponse(iter(), media_type="text/plain")
几个实践细节:
StreamingResponse本身接受同步生成器:FastAPI/Starlette 会在线程池中迭代它,因此这里不需要手动包 executor——同步迭代器可以直接被框架消费,这正是文档示例的写法;media_type按需选择:text/plain是最简单直观的选择;如果要在前端用EventSource或fetch流式读取,可以改为text/event-stream并在生成器内自行拼装 SSE 帧;Runner应复用:Runner构造时会加载权重,应在进程生命周期内只创建一次,多个请求共享(getting-started 指南同样强调了这一点),避免每次请求都重新加载模型导致巨大的启动开销。
生产环境:优先用 HTTP 服务而不是进程内 Runner
原文档明确给出了一条重要的生产建议:在 Web 应用进程内加载模型仅适合原型验证;生产环境应把 mistralrs 作为独立 HTTP 服务器运行,再用 OpenAI Python 客户端调用它。理由是 HTTP 服务器在高负载下的流式行为更健壮。
具体做法参见 OpenAI-compatible API 服务指南:
mistralrs serve -m Qwen/Qwen3-4B
然后使用 OpenAI SDK 以 http://localhost:1234/v1 为 base URL 发起流式请求,api_key 客户端必填但服务端不校验。stream=True 时即可获得逐 token 输出。这样模型进程与 Web 应用进程隔离,模型崩溃、内存压力不会拖垮 Web 服务,也便于独立扩缩容与监控。
流式过程中的错误处理
流式生成可能在中途失败:内存不足(OOM)、生成失败、校验错误等。SDK 的行为是:迭代器以 ValueError 作为下一个迭代项抛出引擎的错误消息,而不是中断已产生的输出——已 yield 的 chunk 不受影响,部分输出得以保留。
import sys
try:
for chunk in stream:
delta = chunk.choices[0].delta.content
if delta:
print(delta, end="", flush=True)
except ValueError as e:
print(f"\n\nStream ended: {e}", file=sys.stderr)
结合源码(stream.rs),可以给出更精确的行为说明:
- 错误类型:
ValueError(内部是PyValueError),异常消息即引擎返回的错误文本。触发它的引擎事件包括Response::ModelError(生成失败、OOM 等)、Response::ValidationError(请求参数校验失败)和Response::InternalError(内部错误)。 - 部分输出保留:由于错误是"下一个迭代项",循环里已处理的 chunk 已经输出到客户端,不会回滚。这在长回复生成到一半崩溃时尤其重要——用户至少能看到已经生成的内容。
- 通道关闭的兜底:如果接收端返回
None(通道关闭且无更多消息),__next__也会返回ValueError("Received none in ChatCompletionStreamer")。 - 错误处理建议:不要把
ValueError当作"正常结束",正常结束应表现为迭代器耗尽(StopIteration);捕获ValueError后应在响应流中追加错误提示(如示例中打印到stderr),并视业务场景决定是否记录日志、上报指标或触发重试。
迭代器的终止条件
迭代器何时结束?源码中给出了精确判定:当某个 chunk 的所有 choice 都带有 finish_reason 时,is_done 置为 true,下一次 __next__ 返回 None。
Response::Chunk(response) => {
if response.choices.iter().all(|x| x.finish_reason.is_some()) {
this.is_done = true;
}
return Some(Ok(response));
}
也就是说,终止发生在"携带 finish_reason 的最后一个 chunk 被 yield 之后"。此时该 chunk 的 delta.content 通常为 None(这也是示例中 if delta: 判断存在的原因),随后循环自然退出。若某个 choice 有 finish_reason 而其他 choice 还没有,迭代仍会继续,直到全部 choice 都结束。
服务端工具与流式输出:透明跳过进度事件
当生成过程中运行了服务端工具(Web 搜索、代码执行、shell、MCP 工具等)时,引擎会向通道发送工具进度事件。如前所述,ChatCompletionStreamer 在内部循环中通过 continue 把这些事件全部跳过,Python 侧只会收到内容 chunk,不会看到任何工具执行的中间状态:
Response::AgenticToolCallProgress { .. } => continue,
Response::AgenticToolApprovalRequired { .. } => continue,
Response::BlockDenoisingProgress(_) => continue,
Response::File(_) => continue,
对于绝大多数场景——只需要把生成的文本流式展示给用户——这种透明跳过正是期望行为:无需任何额外代码,工具调用对流式消费者完全不可见。
但如果你的应用需要观察工具执行进度(例如展示"正在搜索…"、"正在执行代码…"),则应改用 agentic runtime 的事件流。该指南明确指出:目前最完整的面向应用的事件流是 /v1/chat/completions 加 stream: true,它会输出标准 OpenAI 兼容 chunk,外加 mistral.rs 的 agentic_tool_call_progress Server-Sent Events(SSE)事件,覆盖模型输出、服务端工具执行、生成媒体、文件与持久会话状态等运行时环节。
小结
mistral.rs Python SDK 的流式输出由 ChatCompletionStreamer 驱动,其核心契约可总结为四点:
- 同步迭代:
stream=True返回同步迭代器,异步代码需用loop.run_in_executor(None, next, stream, None)包装; - Web 集成:FastAPI 的
StreamingResponse可直接消费同步生成器;生产环境建议改用 mistralrs serve + OpenAI 客户端; - 错误即迭代项:流中途失败以
ValueError形式作为下一个迭代项返回,部分输出保留,可用try/except ValueError捕获; - 透明工具事件:服务端工具进度事件被内部跳过,纯文本流式无感知;需要进度事件请走 agentic runtime 的 SSE 流。
这些行为均可从 mistralrs-pyo3 源码、类型声明 与 可运行示例 中直接验证,适合作为接入流式聊天、SSE 推送或 Agent 应用时的可靠依据。
更多推荐

所有评论(0)