流式处理 (Streaming)
🌊 LangChain 中实现流式输出的核心方法是什么?stream 和 astream 的区别。¶
在 LangChain 中,流式输出是指模型每生成一个 token 就立即返回给客户端,而不是等整个回答生成完。实现流式输出的核心方法是 stream(同步)和 astream(异步),它们是 Runnable 接口的一部分,所有 LCEL 链、模型、组件都支持。
-
stream(input, config):返回一个 Python 迭代器,遍历它能逐个获取事件。对于 LLM,每个事件通常是一个AIMessageChunk,携带增量内容(token)。它阻塞当前线程,适合同步环境。 -
astream(input, config):返回一个异步迭代器,适用于asyncio环境。它不阻塞事件循环,可以与其他异步任务并发执行,是构建高性能 Web 服务的基石。
两者本质上都是生产-消费模型的封装,底层依赖 LLM 自身的流式 API(例如 OpenAI 的 stream=True),LangChain 内部通过回调系统(CallbackManager)把这些 token 转化为事件流。
区别:
在实际项目中,如果你使用 FastAPI 等异步框架,必须使用 astream,否则会阻塞整个服务;而在 Jupyter Notebook 或命令行脚本中,stream 更为便捷。
🛠️ 在自定义链中,如何实现流式输出?需要继承哪些方法?¶
如果你要编写一个自定义链并支持流式输出,需要遵循 LangChain 的 Runnable 协议。相比于传统的 Chain 基类,LCEL 的 Runnable 是更推荐的实现方式。
关键步骤:
-
继承
Runnable而不是Chain。Runnable提供了统一的接口,包括invoke、stream、astream等。 -
实现核心方法:
invoke(self, input, config):返回完整结果,用于非流式调用。stream(self, input, config):返回同步迭代器。通常你需要重写此方法,在内部调用下游 LLM 的stream,并处理中间逻辑,然后yield每个 chunk。transform(self, input, config):这是一个生成器函数,框架默认的stream就是调用它。你可以重写transform来逐个产出事件。-
对于异步,还有
ainvoke、astream、atransform。 -
使用
@chain装饰器(LangChain 提供)可以快速将一个函数转换为支持流式的Runnable,但自定义复杂链时手动继承更灵活。
示例:一个简单的自定义流式链
from langchain_core.runnables import Runnable
from typing import Iterator
class MyStreamingChain(Runnable):
def __init__(self, llm):
self.llm = llm
def invoke(self, input, config=None):
return self.llm.invoke(input)
def stream(self, input, config=None) -> Iterator:
# 直接委托给 LLM 的 stream,返回每个 token chunk
for chunk in self.llm.stream(input):
yield chunk
如果你想在中间做一些处理,比如在生成 token 后添加一个前缀,可以在 stream 中实现:
def stream(self, input, config=None):
yield "🤖 " # 自定义输出
for chunk in self.llm.stream(input):
yield chunk
重要提示:要确保内部组件(如 LLM)已经开启流式模式,否则 stream 返回的可能是一次性结果而不是迭代器。
📝 写出一个简单的 LCEL 链,并使用 stream 方法逐 token 打印结果。¶
下面是一个经典的 LCEL 链:Prompt → LLM → StrOutputParser,然后流式输出。
from langchain_core.prompts import ChatPromptTemplate
from langchain_core.output_parsers import StrOutputParser
from langchain_openai import ChatOpenAI
# 1. 创建组件
prompt = ChatPromptTemplate.from_template("请用50字以内解释:{topic}")
llm = ChatOpenAI(model="gpt-3.5-turbo", streaming=True) # 必须开启 streaming
output_parser = StrOutputParser()
# 2. LCEL 链
chain = prompt | llm | output_parser
# 3. 流式调用
topic = "量子计算"
print("🤖 ", end="", flush=True)
for chunk in chain.stream({"topic": topic}):
print(chunk, end="", flush=True) # chunk 就是一个个字符串
print() # 换行
执行效果:控制台会像打字机一样逐个字符出现最终答案。
关键点:
-
LLM 必须设置
streaming=True,否则stream会退化为一次返回全部内容。 -
StrOutputParser在流式模式下能够处理字符串 chunk,并且会自动拼接,但输出的 chunk 仍然是增量文本。 -
如果链中没有
StrOutputParser,stream返回的将是AIMessageChunk对象,需要手动提取.content。
❓ 流式输出时,LangChain 的 Output Parser 还能正常工作吗?如何处理不完整 token?¶
可以,但有限制。 有些 Output Parser 需要看到完整的生成文本才能解析(例如 JSON 解析器),而流式模式下文本是逐步到达的,这就会出问题。
-
简单 Parser(如
StrOutputParser、PydanticOutputParser的流式子类):能够处理增量输入。StrOutputParser在流式模式下会直接透传字符串 chunk,并在最终调用时拼接,所以没有问题。 -
结构化 Parser(如
JsonOutputParser、PydanticOutputParser):默认不支持流式,因为它们需要完整 JSON 才能反序列化。强行在流式模式下使用,可能会在接收到不完整 JSON 时抛出解析错误。
如何处理不完整 token?
-
使用支持流式的 Parser:LangChain 提供了一些专为流式设计的输出解析器,例如
JsonOutputToolsParser、PydanticOutputParser的流式变体(部分存在)。你可以在langchain.output_parsers中查找支持stream的解析器。 -
缓冲直到完整:在自定义流式处理中,收集所有 chunk,直到生成结束,再一次性解析。这牺牲了流式的即时性,但保证了解析正确。
-
增量解析器:对于一些格式(如 JSON Lines),可以实现一个状态机,每收到一行就解析一次。但 LangChain 目前没有内置通用的增量解析器。
-
后处理:在不改变 Parser 的情况下,前端可以累积字符,当检测到完整消息(例如 JSON 闭合、特定分隔符)后再渲染。但这样 Parsing 的负担转移到了客户端。
实践建议:对于需要结构化输出的场景,尽量让 LLM 返回简单文本,然后通过后续的链步骤或工具调用进行结构化,而不是依赖流式输出时的 Parser。如果必须流式输出 JSON,推荐使用 JsonOutputToolsParser,因为它基于工具调用的 JSON,天然适合流式。
🧠 在 Agent 中使用流式输出时,你会看到什么?Agent 的 Thought 也会流式吗?¶
Agent 的流式输出与普通链有所不同,因为 Agent 内部是多步推理的。你会看到每一步的中间结果,而不仅仅是最终答案。
默认行为(使用 AgentExecutor.stream):
-
每个完整步骤结束后,会产生一个事件,包含该步骤的
Action或Final Answer。这是经典的“步级流式”,而不是 token 级流式。 -
例如,Agent 先思考、调用工具,你会收到一个事件
{"actions": [AgentAction(...)], ...};工具执行完,产生{"steps": [...], ...};最后输出{"output": "最终答案", ...}。
是否能看到 token 级的 Thought 流式?
-
标准 ReAct Agent 不支持:ReAct Agent 的 Prompt 要求 LLM 一次性输出
Thought, Action, Action Input,因此 Thought 不会流式出现,而是整体生成完才作为一步被解析。 -
OpenAI Functions Agent 也不支持 token 级 Thought 流式,因为其 Thought 蕴含在函数调用中。
-
如果你想看到逐 token 的思考过程,可以:
- 使用
StreamingStdOutCallbackHandler,但那是把所有 LLM 输出(包括 Thought)都流式打印,但不区分是否是 Thought。 - 自定义一个 Agent,将思考过程(Thought)作为独立的流式块产出,但需要重构输出解析器。
目前在 LangChain 中,Agent 的流式输出主要是步级流式,这对于大多数应用已经足够,因为用户更关心最终答案的逐步生成,而不是内部的推理碎片。 如果你需要展示“思考中…”的动画,可以在 on_agent_action 回调中更新 UI。
⌨️ 如何实现一个“打字机效果”的控制台输出?在 LangChain 中可以很方便做到吗?¶
很方便。 LangChain 提供了 StreamingStdOutCallbackHandler,只需将 LLM 设置为 streaming=True 并传入该回调即可。
示例:
from langchain_openai import ChatOpenAI
from langchain.callbacks.streaming_stdout import StreamingStdOutCallbackHandler
llm = ChatOpenAI(
model="gpt-3.5-turbo",
streaming=True,
callbacks=[StreamingStdOutCallbackHandler()]
)
chain = prompt | llm | StrOutputParser()
chain.invoke({"topic": "黑洞"})
运行时,控制台会逐字输出结果,没有额外的编程。
自定义“打字机效果”:
如果你需要更复杂的控制(例如添加光标闪烁、颜色),可以自己实现一个回调:
from langchain.callbacks import BaseCallbackHandler
import sys, time
class TypewriterHandler(BaseCallbackHandler):
def on_llm_new_token(self, token: str, **kwargs) -> None:
sys.stdout.write(token)
sys.stdout.flush()
time.sleep(0.02) # 模拟人类打字速度
将其实例传入 LLM 的 callbacks 即可。注意:添加延迟会显著拖慢生成速度,仅适合演示,生产环境不宜使用。
LCEL 的 stream 方法也可以直接实现打字机效果:
这种方法不依赖回调,更灵活。
🔍 如果你想在流式输出过程中,基于已生成的内容做一些实时处理(如敏感词过滤),怎么做?¶
可以利用 stream 迭代器或自定义回调,在 token 到达时立即检查。
方案一:在消费流时过滤(推荐)
blocked_words = ["暴力", "色情"]
for chunk in chain.stream({"topic": topic}):
if any(word in chunk for word in blocked_words):
# 检测到敏感词,停止输出并返回警告
print("[内容已屏蔽]")
break
print(chunk, end="", flush=True)
这种方式简单,但过滤粒度较粗(token 可能跨词切分,导致漏检)。
方案二:使用回调进行缓冲过滤
实现一个回调处理器,在 on_llm_new_token 中维护一个缓冲区,按逻辑边界(如遇到标点、空格、换行)检查缓冲区内容。如果检测到敏感词,可以抛出异常中断生成,或者替换后再输出。
class ContentFilterHandler(BaseCallbackHandler):
def __init__(self, blocked_words):
self.buffer = ""
self.blocked = set(blocked_words)
self.cancelled = False
def on_llm_new_token(self, token: str, **kwargs) -> None:
if self.cancelled:
return
self.buffer += token
# 检测完整词(简单策略:遇到空格或标点)
if token in (' ', '\n', '.', ',', '。'):
for word in self.blocked:
if word in self.buffer:
self.cancelled = True
raise StopIteration("敏感词中断")
sys.stdout.write(self.buffer)
self.buffer = ""
在生产环境,建议使用专门的内容安全服务(如 OpenAI Moderation API),通过异步调用避免阻塞主线程。
🌐 在 Web 应用中,LangChain 的流式输出如何通过 Server-Sent Events (SSE) 发送到前端?¶
SSE 是服务端向客户端推送事件的标准协议。在 Python Web 框架(如 FastAPI)中,可以结合 LangChain 的异步流式轻松实现。
FastAPI + SSE 示例:
from fastapi import FastAPI
from fastapi.responses import StreamingResponse
from langchain_openai import ChatOpenAI
from langchain_core.prompts import ChatPromptTemplate
from langchain_core.output_parsers import StrOutputParser
import asyncio, json
app = FastAPI()
llm = ChatOpenAI(model="gpt-3.5-turbo", streaming=True)
prompt = ChatPromptTemplate.from_template("请解释:{topic}")
chain = prompt | llm | StrOutputParser()
async def event_generator(topic: str):
async for chunk in chain.astream({"topic": topic}):
# SSE 格式:data: {json}\n\n
yield f"data: {json.dumps({'token': chunk})}\n\n"
await asyncio.sleep(0) # 让出控制权
yield "data: [DONE]\n\n"
@app.get("/stream")
async def stream(topic: str):
return StreamingResponse(event_generator(topic), media_type="text/event-stream")
前端使用 EventSource API 接收:
const eventSource = new EventSource("/stream?topic=量子计算");
eventSource.onmessage = (event) => {
if (event.data === "[DONE]") {
eventSource.close();
return;
}
const data = JSON.parse(event.data);
document.getElementById("output").innerText += data.token;
};
关键点:
-
必须使用
astream(异步),否则会阻塞 FastAPI 事件循环。 -
SSE 的响应头
Content-Type: text/event-stream由StreamingResponse自动设置。 -
客户端连接中断时,生成器会抛出
GeneratorExit,可以在finally中清理资源。
🚀 LangServe 是如何支持流式输出的?客户端如何消费?¶
LangServe 是 LangChain 的部署工具,能自动将链发布为 REST API。它原生支持流式输出,只需在链配置中启用。
服务端:使用 add_routes 注册链,LangServe 会自动为支持流式的链生成 /stream 端点。
from langserve import add_routes
from fastapi import FastAPI
app = FastAPI()
# 假设 chain 支持 stream(LLM 开启了 streaming)
add_routes(app, chain, path="/chat")
LangServe 会暴露 POST /chat/stream 接口,接收 JSON 输入,返回 SSE 流。
客户端消费:LangServe 提供了 RemoteRunnable,可以像本地调用一样处理流式。
from langserve import RemoteRunnable
remote_chain = RemoteRunnable("http://localhost:8000/chat/")
for chunk in remote_chain.stream({"topic": "AI"}):
print(chunk)
你也可以直接用 requests 库获取 SSE 流,解析 data: 行。
特点:
-
LangServe 自动处理了序列化、错误传播、配置传递。
-
流式接口返回的是标准的 SSE 事件流,内容为 LangChain 的内部事件(如
AIMessageChunk),客户端需要解析。 -
LangServe 的 playground (Swagger UI) 也支持流式测试。
🛑 当流式输出被中断(如用户关闭了网页),后端如何感知并中止生成?¶
在 Web 应用中,当客户端断开 SSE 连接时,服务端的生成器会收到一个 GeneratorExit 异常(或 asyncio.CancelledError 对于 astream)。我们可以捕获这个异常,并立即中止 LLM 调用,释放资源。
异步 FastAPI 示例:
async def event_generator(topic: str):
try:
async for chunk in chain.astream({"topic": topic}):
yield f"data: {json.dumps({'token': chunk})}\n\n"
except asyncio.CancelledError:
# 客户端断开连接,asyncio 任务被取消
print("客户端已断开,生成中止")
# 这里可以调用 LLM 的取消方法(如果支持)
# 或者简单地清理资源
return
finally:
print("清理资源")
对于同步 WSGI 应用(如 Flask):
def generate():
for chunk in chain.stream(input):
try:
yield f"data: {chunk}\n\n"
except GeneratorExit:
# 客户端断开
break
底层 LLM 的中止:目前大多数 LLM API(如 OpenAI)不支持从客户端主动中止正在进行的流式请求。连接断开后,API 通常还会继续生成,但结果会被丢弃。要真正中止,需要依赖 HTTP/2 的 RST_STREAM 或 WebSocket 的关闭帧,LangChain 未对此做封装。在实际应用中,通常可以忽略这部分开销,因为断开连接已经阻止了数据继续传输,GPU 资源会在后台自动释放。
最佳实践:
-
在
finally块中记录日志,统计生成中断率。 -
如果使用的是本地模型,可以尝试发送中断信号给推理进程(需要自定义集成)。
-
对于成本敏感的场景,可以设置一个“最大生成时间”,超时自动取消任务,避免无效生成。
⚡ 在异步环境下,astream 如何配合 FastAPI 的 StreamingResponse 使用?¶
将 LangChain 的异步流式与 FastAPI 结合,是构建高性能 LLM 服务的关键。核心在于 astream 返回的是一个异步迭代器,而 FastAPI 的 StreamingResponse 可以接收一个异步生成器,实现服务端推送。
实现步骤:
-
确保 LLM 开启
streaming=True,并使用异步客户端(如AsyncChatOpenAI)。 -
定义异步生成器函数:在函数内部使用
async for遍历chain.astream(input),将每个 chunk 格式化为 SSE 消息,然后yield。 -
路由返回
StreamingResponse,传入该生成器,并指定media_type="text/event-stream"。
示例代码:
from fastapi import FastAPI
from fastapi.responses import StreamingResponse
from langchain_openai import ChatOpenAI
from langchain_core.prompts import ChatPromptTemplate
from langchain_core.output_parsers import StrOutputParser
import json, asyncio
app = FastAPI()
llm = ChatOpenAI(model="gpt-3.5-turbo", streaming=True) # 异步自动适配
prompt = ChatPromptTemplate.from_template("请解释:{topic}")
chain = prompt | llm | StrOutputParser()
async def event_generator(topic: str):
try:
async for chunk in chain.astream({"topic": topic}):
# 每个 chunk 是字符串
data = json.dumps({"token": chunk}, ensure_ascii=False)
yield f"data: {data}\n\n"
await asyncio.sleep(0) # 让出控制权,避免阻塞事件循环
except asyncio.CancelledError:
# 客户端断开连接
pass
finally:
yield "data: [DONE]\n\n"
@app.get("/stream")
async def stream_endpoint(topic: str):
return StreamingResponse(event_generator(topic), media_type="text/event-stream")
关键点:
-
必须使用
async for,否则StreamingResponse无法异步消费。 -
await asyncio.sleep(0)主动让出 CPU,保证高并发下的公平性,但不是必须。 -
捕获
asyncio.CancelledError处理客户端断开,避免资源泄露。 -
生成器函数
event_generator是异步的,FastAPI 会在后台通过事件循环驱动它。
踩坑提醒:如果 LLM 使用了同步的 stream 而不是 astream,FastAPI 会在线程池中执行,失去了异步的优势,高并发时性能会急剧下降。务必使用 AsyncChatOpenAI 或 ChatOpenAI 默认的异步支持。
⏱️ 流式输出中,如果 LLM 返回速度很快,会不会导致 UI 更新过于频繁?如何限流?¶
会的。如果模型生成速度极快(例如每秒 100+ token),前端 DOM 更新过于频繁会导致页面卡顿、CPU 飙升,用户体验反而变差。此时需要在前端或后端进行节流。
策略一:后端节流(推荐)
在 LangChain 的流式回调或生成器中进行批量输出。例如,每累积 N 个 token 或者每隔一定时间发送一次。
import time
async def throttled_generator(chain, input, min_interval=0.05):
buffer = ""
last_send = time.time()
async for chunk in chain.astream(input):
buffer += chunk
now = time.time()
if now - last_send >= min_interval:
yield buffer
buffer = ""
last_send = now
if buffer:
yield buffer
这种方式简单有效,直接减少网络包数量。但会导致前端的“打字机”效果略显卡顿,需要权衡。
策略二:前端节流
前端使用 requestAnimationFrame 或定时器(如 50ms)合并渲染。这是最灵活的方式,后端无需改动。
let buffer = "";
const render = () => {
if (buffer) {
document.getElementById("output").innerText += buffer;
buffer = "";
}
requestAnimationFrame(render);
};
requestAnimationFrame(render);
eventSource.onmessage = (e) => {
const data = JSON.parse(e.data);
buffer += data.token;
};
策略三:利用 on_llm_new_token 回调进行节流
实现一个自定义回调,在 on_llm_new_token 中维护缓冲区,并结合定时器发送。
class ThrottleHandler(BaseCallbackHandler):
def __init__(self, queue, interval=0.05):
self.queue = queue
self.interval = interval
self.buffer = ""
self.last_send = time.time()
def on_llm_new_token(self, token, **kwargs):
self.buffer += token
if time.time() - self.last_send >= self.interval:
self.queue.put_nowait(self.buffer)
self.buffer = ""
self.last_send = time.time()
结论:通常前端节流最简单,后端节流可节省带宽。对于高并发服务,后端节流能显著降低系统负载。
📊 如何统计流式输出的总 token 数?回调在这里扮演什么角色?¶
流式模式下,token 是逐个到达的,因此无法像非流式那样直接从返回对象中拿到 token_usage。回调则是捕获这些 token 并进行累加的理想场所。
方法一:利用 on_llm_new_token 回调计数
每次回调触发时,将计数器加一。但这只能统计生成的 token 数(completion tokens),无法直接获取 prompt tokens。要获得完整统计,需要在 on_llm_end 中读取 LLM 返回的 usage 信息(OpenAI 等流式 API 在最后一条 chunk 中会包含 usage)。
class TokenCounter(BaseCallbackHandler):
def __init__(self):
self.prompt_tokens = 0
self.completion_tokens = 0
self.total_tokens = 0
def on_llm_new_token(self, token: str, **kwargs):
self.completion_tokens += 1
def on_llm_end(self, response, **kwargs):
# 尝试从 response 中获取 usage(如果流式最后携带)
if hasattr(response, 'llm_output') and response.llm_output:
usage = response.llm_output.get('token_usage', {})
self.prompt_tokens = usage.get('prompt_tokens', 0)
self.completion_tokens = usage.get('completion_tokens', self.completion_tokens) # 修正
self.total_tokens = usage.get('total_tokens', self.prompt_tokens + self.completion_tokens)
else:
self.total_tokens = self.prompt_tokens + self.completion_tokens
方法二:使用 stream 迭代器自行统计
如果你在消费 stream 迭代器,可以边生成边计数,但 prompt tokens 无法获得。只能通过预先用 tokenizer 计算 prompt 长度来估算。
import tiktoken
enc = tiktoken.encoding_for_model("gpt-3.5-turbo")
prompt_tokens = len(enc.encode(prompt))
completion_tokens = 0
for chunk in chain.stream(...):
completion_tokens += 1
回调角色:回调是横切关注点的利器,它让你无需修改业务代码就能全局统计 token,配合 run_id 还能区分不同的 LLM 调用。
❌ 如果在流式输出过程中发生错误,你会怎么通知前端?¶
流式场景下的错误处理要保证用户能感知到异常,且不会丢失已生成的内容。关键在于在 SSE 流中发送错误事件,以及前端对异常事件的监听。
后端实现:
在生成器内部捕获异常,然后发送一个特殊的 SSE 事件(例如 {"error": "发生错误,请稍后重试"}),并关闭连接。
async def event_generator(topic: str):
try:
async for chunk in chain.astream({"topic": topic}):
yield f"data: {json.dumps({'token': chunk})}\n\n"
except Exception as e:
error_msg = f"生成失败:{str(e)}"
yield f"data: {json.dumps({'error': error_msg})}\n\n"
finally:
yield "data: [DONE]\n\n"
前端处理:
eventSource.onmessage = (event) => {
if (event.data === "[DONE]") {
eventSource.close();
return;
}
const payload = JSON.parse(event.data);
if (payload.error) {
displayError(payload.error);
eventSource.close();
return;
}
if (payload.token) {
outputDiv.innerText += payload.token;
}
};
注意事项:
-
网络中断或超时不属于应用层异常,前端需要监听
EventSource.onerror并尝试重连。 -
已经显示的部分内容建议保留,让用户知道生成中断了。
-
如果 LLM 支持,可以通过回调提前终止生成(如 OpenAI 的
cancel),但 LangChain 未直接暴露,需要借助底层客户端。
🔀 你如何处理同时有多个流式请求的情况?如何保证每个请求的 stream 不被混淆?¶
LangChain 的流式迭代器本身是请求隔离的。每个 astream 调用返回独立的异步生成器,内部状态(如 LLM 的 HTTP 连接、缓冲区)都是独立的,因此不会混淆。
在 Web 应用中,每个客户端请求都会创建独立的生成器实例,FastAPI 会为每个请求分配独立的协程,流自然不会串扰。
但要注意:
-
共享状态:如果在链或工具中使用了全局变量或类级别的状态,多个请求会竞争。务必使用请求级上下文,或者通过
config传递用户相关信息。 -
并发限制:大量并发流式请求会耗尽 LLM API 的速率限制或系统连接池。需要设置并发控制,例如使用
asyncio.Semaphore。
semaphore = asyncio.Semaphore(10) # 最大并发10
async def event_generator(topic: str):
async with semaphore:
async for chunk in chain.astream({"topic": topic}):
yield ...
- 内存泄漏:如果生成器没有被正确消费(如客户端断开),且后端未做清理,可能会导致缓冲区堆积。务必在
finally中释放资源。
结论:LangChain 的流式设计天然支持多请求隔离,开发者需要关注的是资源管理和并发控制。
🧩 在 LangChain 中,Non-LLM 组件(如 Retriever)是否也支持流式?通常不需要,为什么?¶
LangChain 的 Runnable 接口理论上支持所有组件实现 stream,但Retriever 等非生成式组件通常不实现流式,因为:
-
语义不符:流式意味着逐步产出数据,而检索是一次性操作:输入查询,返回一组完整文档。没有中间状态。
-
性能无益:检索结果通常体量小(几个到几十个文档),一次性返回远比逐文档流式更高效。
-
框架设计:LCEL 的
stream对于非 LLM 组件,通常会退化为一次invoke后包装成单元素的迭代器,不会真正流式。 -
下游需求:即使 Retriever 流式产生文档,下游 LLM 也需要等待所有文档才能生成,因此没有意义。
不过,某些特殊情况可以实现“伪流式”:比如将大量文档分批返回,或者对检索结果进行逐步过滤,但这通常不是标准实践。如果你需要“逐步输出检索结果”,可以自定义 Runnable,但实用价值有限。
📜 流式输出的日志记录一般怎么做?每个 token 都记录吗?¶
通常不记录每个 token,因为日志量太大,且 I/O 开销会严重影响生成性能。
最佳实践:
-
只记录请求和响应摘要:在
on_llm_start记录 prompt 摘要,on_llm_end记录完整回复文本、token 数、耗时。这足以满足审计和监控。 -
采样记录:在生产环境,可对 1% 的请求记录完整回复文本,用于质量评估。
-
错误全记录:流式生成出现异常时,记录完整上下文。
-
使用异步日志:避免日志写入阻塞主线程,可使用
logging.handlers.QueueHandler或异步日志库。
示例:
class StreamLogger(BaseCallbackHandler):
def on_llm_end(self, response, **kwargs):
text = "".join([g[0].text for g in response.generations]) if response.generations else ""
logging.info(f"LLM call finished: tokens={...}, response={text[:200]}...")
为何不记录每个 token:假设每秒生成 50 token,1k 并发,日志写入会成为瓶颈,且 token 片段无完整语义,分析价值低。
🔄 流式处理和批处理的性能有什么不同?什么时候该用 stream,什么时候该用 invoke?¶
性能差异:
-
首 token 延迟:流式极低(毫秒级),用户立即看到反馈;批处理(invoke)需要等待完整生成,首 token 延迟等于总生成时间。
-
总延迟:两者相近,流式可能因网络和调度略高。
-
吞吐量:批处理可以合并请求批量推理,在高并发时吞吐量更高;流式通常每个请求占用一个连接,服务端资源消耗较大。
-
内存占用:流式无需在服务端缓存全部回复,内存友好;批处理需要缓存完整输出。
选择场景:
-
用户交互(聊天、写作助手):必须用 stream,提升感知速度。
-
批量处理(数据标注、报告生成):用 invoke,服务端可以批量调用,成本更低。
-
后端服务间调用:用 invoke,因为调用方不需要实时展示。
-
需要结构化输出的场景:用 invoke,等完整结果后解析,流式解析容易出错。
混合策略:对用户端流式输出,服务端内部对下游 LLM 仍使用流式,但可以通过缓冲将 token 组装后以完整 JSON 返回给内部服务。
❓ 如果上游 LLM 不支持流式输出,LangChain 能模拟流式吗?¶
不能真正模拟 token 级的流式,但可以模拟一种“伪流式”效果。
LangChain 本身没有内置的模拟流式功能。如果 LLM 不支持原生流式,stream 方法会退化为 invoke,一次性返回结果,然后 LangChain 把这个结果包装成一个单元素的迭代器。前端会在一瞬间收到全部文本,失去了打字机效果。
变通方案:
-
客户端模拟:后端返回完整文本,前端通过 JavaScript 定时逐字渲染,实现假流式。这能改善感知延迟,但首字等待时间不变。
-
服务端分块发送:在自定义链中,拿到完整结果后,按句子或固定长度切片,通过生成器逐步 yield,但这并不是真正的流式,且增加了服务端开销。
-
换用支持流式的模型:例如从
text-davinci-003切换到gpt-3.5-turbo(支持 streaming)。这是根本解决之道。
总之,LangChain 无法让不支持流式的 API 产生流式输出,但可以在应用层包装出类似效果,核心瓶颈是首 token 延迟无法消除。
🌊 谈谈你对流式输出中 buffer 和背压控制的理解。¶
流式输出中,buffer(缓冲区) 和 背压(backpressure) 是保证系统稳定性的重要机制。
Buffer 的作用:
-
在生产者(LLM)和消费者(网络传输或前端)之间暂存数据,平衡速度差异。
-
例如,后端使用
asyncio.Queue作为缓冲,生成器向队列放 token,SSE 发送协程从队列取数据。这样即使网络偶尔抖动,生成也不会被阻塞。 -
但 buffer 不能无限大,否则会导致内存溢出。
背压控制:
-
当消费者(客户端)处理速度跟不上生产者(LLM)时,需要让生产者慢下来。
-
在 TCP 层面,网络协议栈会自动施加背压:如果接收窗口满,发送方被阻塞。
-
在应用层,如果我们自定义了队列,可以设置
maxsize,当队列满时,put操作会阻塞或引发异常,从而将压力反向传导到 LLM 生成(迫使生成暂停)。 -
在 LangChain 中,如果使用
astream并直接在async for循环中yield,没有显式 buffer,这时背压依赖于 Web 框架和 ASGI 服务器的内部机制(如 Uvicorn 会管理 TCP 缓冲,当客户端读取慢时,send协程会暂停,进而使生成器协程暂停,达到自然背压)。
最佳实践:
-
使用有界队列(
asyncio.Queue(maxsize=100))作为中间缓冲,避免内存无限增长。 -
监控队列深度,如果持续满载,说明消费能力不足,需要扩容或限流。
-
在生成器内部检测队列状态,如果队列满,适当引入
asyncio.sleep或主动降速(比如降低生成优先级),但通常 Web 服务器和 OS 的背压已足够。
踩坑经验:曾经我在生成大段文本时未设置 buffer,导致客户端网络慢时,生成器一直阻塞在 yield,GPU 空闲,而 LLM API 仍在计时计费。后来引入了一个小 buffer(10个token),并设置超时,如果消费者在超时内未消费,则主动中止生成,避免浪费。这是流式系统中必须考虑的工程细节。