LangChain AI应用开发框架的使用(4) - 聊天模型的流式传输,异步流式输出,深度探索流式传输,使用 LangSmith 跟踪 LLM 应用

发布时间:2026/9/4 8:12:17
LangChain AI应用开发框架的使用(4) - 聊天模型的流式传输,异步流式输出,深度探索流式传输,使用 LangSmith 跟踪 LLM 应用 目录一、聊天模型 -- 流式传输stream() 同步传输astream() 异步传输异步相关概念什么是协程什么是事件循环使用二、使用 StrOutputParser 解析模型的输出三、自定义流式输出解析器四、深度探索流式传输SSE 协议介绍核心特点数据格式LangChain 流式传输流程分析通过源码分析流程LangChain 请求 OpenAI 使用什么协议LangChain 如何支持流式传输OpenAI 返回的块是什么格式如何转换成 AIMessageChunk五、使用 LangSmith 跟踪 LLM 应用一、聊天模型 -- 流式传输流式处理对于使基于 LLM 的应用程序能够响应最终用户至关重要。它通过逐步显示输出甚至在完整的响应准备就绪之前流式传输可以显著改善用户体验。我们之前直接使用 invoke 的调用方式属于非流式传输看到的现象是聊天模型直接返回全量内容若模型思考时间较长则我们等待的时间就越长。伪代码示例我们等待了 20s 之后返回了全量结果。从结果来看我们等待的时间太长了对于用户来说太长的等待时间严重影响体验。而流式传输效果就像 deepseek 客户端那样文字一点点逐段输出不用等全部生成完才展示内容。因此 LangChain 聊天模型原生支持流式返回。stream() 同步传输在 LangChain 聊天模型中可以使用其 .stream() 方法来同步生成流式响应的效果。聊天模型的 .stream() 方法返回一个迭代器该迭代器在生成输出时同步产生输出消息块。可以使用 for 循环实时处理每个块。代码如下打印结果示例通过调试可以看到迭代出来的对象是 AIMessageChunk它代表 AIMessage 的一部分也就是消息块。消息块对象支持直接相加可以把多个 chunk 合并还原完整消息。如下:合并输出结果示例astream() 异步传输对于流式传输通常我们可以选择异步调用。先来了解下异步相关知识。异步相关概念举个场景需要煮一壶水同时还要给朋友发短信。分别用同步传统和异步两种方式完成引入协程和事件循环概念。同步阻塞方式做事必须一件一件来。我们举个例子:总耗时7 秒。问题在 boil_water 函数等待的 5 秒里CPU 完全空闲但却不能去做 send_message 任务效率低下。异步方式使用 asyncio、协程、事件循环。什么是协程多进程通常利用的是多核 CPU 的优势同时执行多个计算任务。每个进程有自己独立的内存管理所以不同进程之间要进行数据通信比较麻烦。多线程是在一个 cpu 上创建多个子任务当某一个子任务休息的时候其他任务接着执行。多线程的控制是由 python 自己控制的。线程存在数据同步问题所以要有锁机制。协程的实现是在一个线程内实现的相当于流水线作业。由于线程切换的消耗比较大所以对于并发编程可以优先使用协程。进程、线程、协程之间的关系协程作为一种轻量级的并发编程模型可以被视为用户态的 “轻量级线程”。与传统线程相比协程的核心优势在于其调度完全由用户空间掌控避免了操作系统内核的频繁介入从而显著降低了上下文切换的开销。在诸如网络数据刷新、资源加载、用户界面更新、以及 I/O 读写等场景下如果并发任务的计算量相对较小、对系统资源占用较低则不必动用操作系统级别的线程。协程的切换则由程序员和编程语言控制程序员决定在何时暂停或恢复协程。协程是一个特殊的函数它可以在执行过程中暂停并在稍后恢复执行。它用 async def 定义并在需要暂停的地方使用 await。在上面这个例子里boil_water_async 和 send_message_async 就是两个协程。什么是事件循环事件循环是 asyncio (Python 标准库中的模块用于编写异步 I/O 操作的代码)的核心你可以把它想象成一个总调度员或一个高效的待办事项 (To‑Do List) 管理员。它的工作流程非常简单它维护着一个任务列表比如煮水、发短信。它不断地循环检查每个任务 a. 如果任务处于 “等待 I/O” 状态比如等水开、等网络响应就暂停它立即去执行下一个已经 “就绪” 的任务。 b. 如果任务的等待时间到了或者 I/O 操作完成了事件循环就恢复执行这个任务。如何运行输出结果总耗时5 秒 (因为两个任务的等待时间是并发的)通过使用 asyncio我们可以在单线程中同时处理多个任务。一个在单线程内调度和管理所有协程的核心机制就是事件循环。它不停地检查哪些协程可以执行哪些在等待。总结一下协程是 asyncio 的核心概念之一。它是一个特殊的函数可以在执行过程中暂停并在稍后恢复执行。协程通过 async def 关键字定义并通过 await 关键字暂停执行等待异步操作完成。要运行一个协程可以使用 asyncio.run() 函数。它会创建一个事件循环并运行指定的协程。事件循环是 asyncio 的核心组件负责调度和执行协程。它不断地检查是否有任务需要执行并在任务完成后调用相应的回调函数。使用可以使用 .astream() 方法来异步生成流式响应的效果这专为非阻塞工作流程而设计。可以在异步代码中使用它来实现相同的实时流式处理行为。代码如下打印结果补充要点.stream()同步流式普通for循环遍历返回的迭代器。.astream()异步流式必须写在async def函数内部使用async for循环来遍历分片。异步函数不能直接调用需要asyncio.run()启动事件循环执行协程。返回的每一个对象同样是 AIMessageChunk 消息块和同步 stream 得到的 chunk 对象类型完全一致。二、使用 StrOutputParser 解析模型的输出还记得最早我们讲过 Runnable 接口Runnable 接口 :聊天模型、输出解析器等组件都实现了 LangChain 的 Runnable 接口他们都是 Runnable 接口的实例。Runnable 定义了一个标准接口允许 Runnable 组件Invoked调用单个输入转换为输出。Batched批处理多个输入被有效地转换为输出。Streamed流式传输输出在生成时进行流式传输。Inspected检查可以访问有关 Runnable 的输入、输出和配置的原理图信息。Composed组合可以组合多个 Runnable以使用 LCEL 协同工作以创建复杂的管道。……可以看到流式传输实际上并不算是聊天模型定义的能力而是只要实现了 Runnable 接口的实例都具备的能力但要注意并非所有组件都必须实现流式处理在某些情况下流式处理要么是不必要的要么很困难要么根本没有意义。例如以后我们会讲解的 Retrievers 检索器就不提供任何流式处理。那么再得出一个关于流式传输的结论 :.stream() 和 .astream() 方法产生的块类型取决于正在流式传输的组件。例如我们当前正在使用聊天模型的流式传输返回的每个块都将是一个 AIMessageChunk。但是对于其他组件块类型可能不同。接下来让我们使用 LCEL 构建一个简单的链该链结合了模型 和 解析器并验证流是否有效。不要忘了使用 LCEL 创建的链也实现了 Runnable 接口。我们将使用 StrOutputParser 来解析模型的输出它从 AIMessageChunk 中提取内容字段为我们提供模型返回的令牌。代码如下打印结果三、自定义流式输出解析器上面我们演示了如何让聊天模型进行流式输出。若此时我们希望修改上一步的输出样式一个字或两个字的输出将输出改为一句话一句话的输出同时保留流式处理功能。那么我们需要在链中使用生成器函数即可完成自定义流式输出的能力。还记得之前说过聊天模型的 .stream() 方法返回的是一个迭代器该迭代器在生成输出时同步产生输出消息块。那么我们的将实现的这些生成器的签名应该是 Iterator[Input] - Iterator[Output]。或者对于异步生成器AsyncIterator[Input] - AsyncIterator[Output]。下面是句号分隔列表的自定义输出解析器的示例打印结果四、深度探索流式传输SSE 协议介绍HTTP 协议本身设计为无状态的请求‑响应模式严格来说是无法做到服务器主动推送消息到客户端但通过 Server‑Sent Events服务器发送事件简称 SSE技术可实现流式传输允许服务器主动向浏览器推送数据流。也就是说服务器向客户端声明接下来要发送的是流消息 (streaming)这时客户端不会关闭连接会一直等待服务器发送过来新的数据流。SSEServer‑Sent Events是一种基于 HTTP 的轻量级实时通信协议浏览器可以通过内置的 EventSource API 接收并处理这些实时事件。核心特点基于 HTTP 协议复用标准 HTTP/HTTPS 协议无需额外端口或协议兼容性好且易于部署。单向通信机制SSE 仅支持服务器向客户端的单向数据推送客户端通过普通 HTTP 请求建立连接后服务器可持续发送数据流但客户端无法通过同一连接向服务器发送数据。自动重连机制支持断线重连连接中断时浏览器会自动尝试重新连接支持 retry 字段指定重连间隔。自定义消息类型客户端发起请求后服务器保持连接开放响应头设置 Content‑Type: text/event‑stream标识为事件流格式持续推送事件流。数据格式服务端向浏览器发送 SSE 数据需要设置必要的 HTTP 头信息每一次发送的消息由若干个 message 组成每个 message 之间由\n\n分隔每个 message 内部由若干行组成每一行都是如下格式Field 可以取值为data [必需]数据内容event [非必需]表示自定义的事件类型默认是 message 事件id [非必需]数据标识符相当于每一条数据的编号retry [非必需]指定浏览器重新发起连接的时间间隔除此之外还可以有冒号:开头的行表示注释。数据示例LangChain 流式传输流程分析LangChain 本身并不 “创造” 或 “规定” 一个底层的网络传输协议而是依赖于其底层的大模型供应商如 OpenAI和我们自身服务应用所使用的 Web 框架如 FastAPI的协议。因此对于 LangChain 的流式传输能力本身是因为大模型供应商提供了流式传输能力由 LangChain 进行调用后接收并处理成一个个的 AIMessageChunk。通过源码分析流程接下来我们将会通过分析相关源码探索整个传输流程。整个过程我们以 OpenAI 举例其他大模型方式类似可自行探索。当我们向 OpenAI 发起流式请求LangChain 实际上会通过 BaseChatOpenAI 类中的 _stream() 方法发起调用。下面来看下 _stream() 方法的关键流程性源码完整源码见 : class langchain_openai.chat_mode l s.base.BaseChatOpenAI:从上述流程看来这就是流式逐块产生 AIMessageChunk 聊天消息的核心方法。那么接下来看三个问题发起调用时底层使用什么协议如何支持流式传输返回的块是什么格式如何转换成 AIMessageChunk这三个问题都掌握后整个流式传输的流程就都能理解了。LangChain 请求 OpenAI 使用什么协议回答这个问题需要看 LangChain 关于 OpenAI 的客户端是怎么定义的。让我们找到class langchain_openai.chat_models._client_utils._SyncHttpxClientWrapper如下所示从上面的代码看来LangChain 使用了 OpenAI 的官方的 OpenAI SDK for Python 接入方式继承了openai._base_client定义了一个 HTTP 客户端。因此在调用时发起的是 HTTP 调用。LangChain 如何支持流式传输开始我们就说了LangChain 本身并不 “创造” 或 “规定” 一个底层的网络传输协议而是依赖于其底层的大模型供应商如 OpenAI的协议。因此当我们发起请求时会在请求中设置 streamTrue (_stream() 源码中的第一步)表示 OpenAI 服务器将在生成 Response 时向客户端发出数据 (server‑sent eventsSSE)。此时 API 会保持 HTTP 连接打开并以特定格式发送数据流。例如我们向原生的 GPT 模型发起一次设置了 streamTrue 的 HTTP 请求你好我是张三。。此时我们会收到来自 OpenAI 的事件块简化后的有效负载序列看了上述示例我们应该可以回答第二个问题。那就是在请求中设置 streamTrue 开启 OpenAI 服务端返回数据块LangChain 通过 _stream() 方法步骤 1、2 完成这件事。OpenAI 返回的块是什么格式如何转换成 AIMessageChunkOpenAI 返回的数据块格式我们已经看到了将其转换为 LangChain 自定义的 AIMessageChunk 则是通过 _convert_chunk_to_generation_chunk() 方法完成的。关键代码如下到此我们就知道了 LangChain 流式传输的完整流程与底层协议。总结一下langchain‑openai 包通过集成 OpenAI Python SDK 提供了一个 HTTP 客户端。因此支持 LangChain 向 OpenAI 的 API 发起调用请求。若希望发起流式传输请求则需在请求中加入 streamTrue向 OpenAI 说明以 SSE 协议进行流式返回。LangChain 接收 OpenAI 的 SSE 格式的响应并将其转换为 LangChain 自封装的消息格式如AIMessageChunk消息。这样就可以以统一的方式处理来自不同模型提供商OpenAI, Anthropic 等的流式响应。五、使用 LangSmith 跟踪 LLM 应用使用 LangChain 构建的许多应用程序可能会包含多个步骤和多次的 LLM 调用。随着这些应用程序变得越来越复杂作为开发者我们能够检查链或代理内部到底发生了什么变得至关重要。最好的方法是使用LangSmith。LangSmith 与框架无关它可以与 langchain 和 langgraph 一起使用也可以不使用。LangSmith 是一个用于帮助我们构建生产级 LLM 应用程序的平台它将密切监控和评估我们的应用。LangSmith 平台地址LangSmith新用户需要注册要想让 LangSmith 跟踪 LLM 应用第一步申请 LangSmith API Key点击 Settings就会跳转到 API Keys 设置页面若没有跳转可以在左侧 tab 栏中找到进入。创建完成后保存好你的 API Key。接下来配置两个环境变量配置完成后让我们任意执行代码查看 LangSmith 平台这将在 LangSmith 的默认跟踪项目中生成调用的跟踪。点击最新一次的调用追踪跟踪会以瀑布流形式展示调用的完整步骤以及每个步骤的详细信息和耗时。让我们能够检查内部到底发生了什么解释RunnableSequence可运行序列就是我们之前讲过的链即我们将model_with_search.invoke() 的结果 (构造成 ToolMessage)当作入参传递给structured_search_model.invoke()。说明本次调用是 RunnableSequence。但不是每次都展示 RunnableSequence根据实际情况而定。ChatOpenAI实际处理的第一步内容调用聊天大模型生成结果。RunnableLambda实际处理的第二步内容表示将 python 可调用对象转换为 Runnable其实就是将 AI 生成的结果转换成为结构化对象。可以看到我们在使用 LangSmith 时没有代码介入只需要配置下环境就可以直接监控我们的应用。