mistral.rs Python SDK 流式输出实战:异步迭代、FastAPI 集成与流中断错误处理

【免费下载链接】mistral.rs Fast, flexible LLM inference 【免费下载链接】mistral.rs 项目地址: https://gitcode.com/GitHub_Trending/mi/mistral.rs

本篇指南围绕 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_penaltytop_ptemperature 等采样参数与 stream=True 的组合使用)。

在写任何进阶用法之前,需要先理解两个关键事实:

  1. SDK 不提供原生异步迭代器send_chat_completion_request 返回的是一个实现了 Python 迭代协议的对象,for chunk in stream 中的 next() 调用是阻塞的。想在 async 代码中使用,必须自行包装(下文详述)。
  2. 每个 chunk 都是 OpenAI 兼容的流式结构choices[0].delta.content 携带一段增量文本,它可能是 None——例如携带 finish_reason 的最后一个 chunk。因此所有示例都先判断 if delta: 再打印。

流式 chunk 的数据模型

mistralrs.pyi 的类型声明可以看到流式响应的完整结构:

  • ChatCompletionChunkResponse:顶层对象,字段包括 idcreatedmodelsystem_fingerprintchoices,以及可选的 usageadapter_generationsession_id
  • ChunkChoice:每个 choice 包含 finish_reasonstr | None)、indexdeltalogprobs
  • Delta:增量内容,contentstr | None)、role、可选的 tool_callsreasoning_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.rsChatCompletionStreamer。它内部持有一个 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)、ValidationErrorInternalError 都被映射为 PyValueError。注意它不是立即抛出,而是作为迭代的"下一个元素"返回,因此已产生的 chunk 不受影响,部分输出得以保留。
  • 工具进度事件被透明跳过AgenticToolCallProgressAgenticToolApprovalRequiredBlockDenoisingProgressFile 等内部事件在循环里 continue 掉,Python 侧只会看到纯内容 chunk。

请求侧的分发逻辑在 mistralrs-pyo3/src/lib.rssend_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 是最简单直观的选择;如果要在前端用 EventSourcefetch 流式读取,可以改为 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/completionsstream: true,它会输出标准 OpenAI 兼容 chunk,外加 mistral.rs 的 agentic_tool_call_progress Server-Sent Events(SSE)事件,覆盖模型输出、服务端工具执行、生成媒体、文件与持久会话状态等运行时环节。

小结

mistral.rs Python SDK 的流式输出由 ChatCompletionStreamer 驱动,其核心契约可总结为四点:

  1. 同步迭代stream=True 返回同步迭代器,异步代码需用 loop.run_in_executor(None, next, stream, None) 包装;
  2. Web 集成:FastAPI 的 StreamingResponse 可直接消费同步生成器;生产环境建议改用 mistralrs serve + OpenAI 客户端;
  3. 错误即迭代项:流中途失败以 ValueError 形式作为下一个迭代项返回,部分输出保留,可用 try/except ValueError 捕获;
  4. 透明工具事件:服务端工具进度事件被内部跳过,纯文本流式无感知;需要进度事件请走 agentic runtime 的 SSE 流。

这些行为均可从 mistralrs-pyo3 源码类型声明可运行示例 中直接验证,适合作为接入流式聊天、SSE 推送或 Agent 应用时的可靠依据。

【免费下载链接】mistral.rs Fast, flexible LLM inference 【免费下载链接】mistral.rs 项目地址: https://gitcode.com/GitHub_Trending/mi/mistral.rs

更多推荐