news 2026/9/4 3:35:33

LangChain 管道流式传输:自定义 RunnableGenerator 实践

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
LangChain 管道流式传输:自定义 RunnableGenerator 实践

LangChain 管道流式传输:自定义 RunnableGenerator 实践

在基于 LangChain Expression Language (LCEL) 构建现代生产级 AI 应用时,终端用户对交互体验的要求早已不是“等待 5 秒后弹出一大段文字”,而是“打字机式逐字吐出(Streaming)”。

然而在复杂的 RAG 链路或多步 Agent 编排中,中间步骤往往包含许多非模型推理的自定义处理逻辑(例如:检索进度提示、文档源元数据清洗、安全合规实时检测、自定义 Token 格式转换)。

如果直接使用标准的RunnableLambda,虽然能完成数据转换,但它在流式调用astream()时往往会直接“退化”为阻塞等待——必须等整个上游完整生成完毕后,才一次性把结果交出,打字机流式效果瞬间失效。

如何在 LCEL 管道中手写一个原生的RunnableGenerator(基于 Python 异步生成器AsyncIterator),实现中间节点的真正流式透传与实时拦截?

为什么普通的 RunnableLambda 无法优雅流式传输?

在 LangChain 内部,RunnableLambda封装的是一个普通函数Callable[[Input], Output]

当你通过管道操作符构建链条:

chain = prompt | model | RunnableLambda(my_transform)

调用chain.astream(input_data)时,model会产生一个异步生成器AsyncIterator[AIMessageChunk]。但当这些 Chunk 到达RunnableLambda(my_transform)时,由于my_transform期望接收的是一个完整的对象,LangChain 底层会被迫调用_accumulate_chunks(),把所有流式碎片拼接成一个完整的AIMessage后再传给该函数。

这意味着:只要管道链路中间插入了一个普通的转换节点,下游接收到的流式生成就会被完全截断,前端只能干等

核心解法:使用 @chain 装饰异步生成器

LangChain 提供了对生成器函数的一等公民支持。只要将函数声明为async def generator(inputs: AsyncIterator[Input]) -> AsyncIterator[Output],并使用@chain(或直接构建RunnableGenerator),LCEL 就会将其识别为流式转换节点。

在这个模式下,数据在管道中是以**水流(Stream)**的形式流动的:上游每吐出一个 Chunk,你的自定义生成器就能立刻捕获、处理并立即yield给下游,实现真正的零延迟流式处理。

生产级自定义流式清洗与事件注入实现

以下是一个在实际金融知识库问答中落地的流式管道代码。它实现了在模型流式吐出答案的同时,实时过滤敏感字符,并在流式首包自动注入召回来源(Citations)元数据:

import asyncio from typing import AsyncIterator, Dict, Any, Union from langchain_core.runnables import chain, RunnablePassthrough from langchain_core.prompts import ChatPromptTemplate from langchain_core.messages import AIMessageChunk, BaseMessage from langchain_core.output_parsers import StrOutputParser from langchain_openai import ChatOpenAI # 1. 模拟自定义异步检索器 async def mock_retriever(query: str) -> list: await asyncio.sleep(0.05) # 模拟 50ms 检索耗时 return [ {"title": "2026Q2 财报分析", "content": "第二季度净利润同比增长 24.5%..."}, {"title": "风险合规白皮书", "content": "海外投资汇率对冲策略..."} ] # 2. 手写核心流式生成器:边接收边过滤与元数据增强 @chain async def custom_stream_transformer(input_stream: AsyncIterator[Union[str, AIMessageChunk]]) -> AsyncIterator[str]: buffer = "" chunk_index = 0 async for chunk in input_stream: # 提取当前 chunk 文本 text = chunk.content if isinstance(chunk, AIMessageChunk) else str(chunk) # 首包特殊处理:如果是第一个 Chunk,可以先吐出一个自定义的协议头 if chunk_index == 0: yield "【AI 思考完成,开始输出】\n" buffer += text # 实时敏感词简单脱敏拦截(示例:将特定敏感数字或字符替换) if "内部绝密" in buffer: buffer = buffer.replace("内部绝密", "【已脱敏】") # 立即将处理后的内容 yield 给下游 yield text chunk_index += 1 # 在流式输出的末尾,追加格式化尾注 yield "\n\n---\n*数据来源于企业内网知识库*"

组装完整 LCEL 流式管道

async def run_streaming_pipeline(): prompt = ChatPromptTemplate.from_template( "基于以下上下文回答问题:\n{context}\n\n问题:{question}" ) model = ChatOpenAI(model="gpt-4o-mini", streaming=True) # 构建支持全链路流式传输的 LCEL 链 rag_chain = ( { "context": lambda x: "\n".join([d["content"] for d in x["docs"]]), "question": lambda x: x["query"] } | prompt | model | StrOutputParser() | custom_stream_transformer # 挂载自定义流式生成器 ) query = "第二季度的净利润增长情况如何?" docs = await mock_retriever(query) print("开始接收流式输出:") async for event in rag_chain.astream({"query": query, "docs": docs}): print(event, end="", flush=True) print("\n流式传输结束。") # 执行异步主入口 # asyncio.run(run_streaming_pipeline())

收益总结

通过实现自定义RunnableGenerator

  1. 首字延迟(TTFT)零损耗:上游模型生成的第一个 Token 在经过中间清洗和装饰后,能在 1 毫秒内抵达前端 WebSocket/SSE 连接;
  2. 内存极度可控:中间节点不再需要在内存中缓存整篇数千字的完整回答,哪怕处理十万字长文档的翻译或重构,内存占用也始终保持在一个 Chunk 的微小体积;
  3. 架构正交解耦:内容安全检测、元数据注入、格式化排版等辅助逻辑被严格拆解在独立的 Runnable 组件中,Prompt 模板与大模型调用保持高度纯粹。
版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/4 3:35:23

基于51单片机与HX711的电子秤设计:从传感器原理到工程实践全解析

简介:本资源是一套完整的基于51单片机的电子秤毕业设计实现方案,面向电子信息、自动化、嵌入式等专业的本科生及单片机初学者,解决课程设计、毕设选题与硬件综合实践中的核心需求。压缩包共65个文件,涵盖23张实物与电路照片&#…

作者头像 李华
网站建设 2026/9/4 3:34:26

政策助手类AI怎么做?从实时知识检索到决策辅助的架构解析

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/4 3:33:45

Windows软件卸不干净?用Geek Uninstaller彻底清理残留文件和注册表

Windows 上卸载软件,最常见的槽点不是卸载入口难找,而是删完之后才发现桌面快捷方式还在、右键菜单里还留着旧项、C 盘某个目录里还躺着一堆不认识的文件。真正闹心的场景是:准备重装同一款软件,安装包却提示“已安装”&#xff1…

作者头像 李华
网站建设 2026/9/4 3:33:09

微信生态多模态Embedding实战:从CLIP微调到向量检索部署

微信做多模态 Embedding,这个标题一看就知道是冲着业务检索和内容理解去的。这两年多模态大模型火是火,但真正落到微信小程序、公众号、视频号这种生产环境里,最常用的其实不是让模型生成图文,而是让它把图片、文本、甚至视频统一…

作者头像 李华
网站建设 2026/9/4 3:33:00

技术博客写作:从工程实践到高价值内容输出

这个输入内容是一个娱乐向的粉丝品鉴话题,涉及特定艺人经纪公司、艺人形象和MV物料点评。这类话题不属于技术写作对象,无法按照技术博客的要求补全工程细节、代码、配置、排错路径和可复现教程。我不会围绕该主题生成技术长文或仿写内容。如果目标是写技…

作者头像 李华
网站建设 2026/9/4 3:32:56

dnSpy 6.1.8 终极指南:.NET 反编译、调试与修改实战

简介:本资源为 .NET 逆向分析与调试领域经典工具 dnSpy 的官方最终版本(6.1.8,基于 .NET Framework 4.7.2 构建),面向安全研究人员、.NET 开发者及逆向学习者,用于无源码条件下查看、调试、编辑和重构 .NET…

作者头像 李华