跳转至

流式处理 (Streaming)

🌊 LangChain 中实现流式输出的核心方法是什么?streamastream 的区别。

在 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 是更推荐的实现方式。

关键步骤:

  1. 继承 Runnable 而不是 ChainRunnable 提供了统一的接口,包括 invokestreamastream 等。

  2. 实现核心方法:

  3. invoke(self, input, config):返回完整结果,用于非流式调用。
  4. stream(self, input, config):返回同步迭代器。通常你需要重写此方法,在内部调用下游 LLM 的 stream,并处理中间逻辑,然后 yield 每个 chunk。
  5. transform(self, input, config):这是一个生成器函数,框架默认的 stream 就是调用它。你可以重写 transform 来逐个产出事件。
  6. 对于异步,还有 ainvokeastreamatransform

  7. 使用 @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 仍然是增量文本。

  • 如果链中没有 StrOutputParserstream 返回的将是 AIMessageChunk 对象,需要手动提取 .content


❓ 流式输出时,LangChain 的 Output Parser 还能正常工作吗?如何处理不完整 token?

可以,但有限制。 有些 Output Parser 需要看到完整的生成文本才能解析(例如 JSON 解析器),而流式模式下文本是逐步到达的,这就会出问题。

  • 简单 Parser(如 StrOutputParserPydanticOutputParser 的流式子类):能够处理增量输入。StrOutputParser 在流式模式下会直接透传字符串 chunk,并在最终调用时拼接,所以没有问题。

  • 结构化 Parser(如 JsonOutputParserPydanticOutputParser):默认不支持流式,因为它们需要完整 JSON 才能反序列化。强行在流式模式下使用,可能会在接收到不完整 JSON 时抛出解析错误。

如何处理不完整 token?

  1. 使用支持流式的 Parser:LangChain 提供了一些专为流式设计的输出解析器,例如 JsonOutputToolsParserPydanticOutputParser 的流式变体(部分存在)。你可以在 langchain.output_parsers 中查找支持 stream 的解析器。

  2. 缓冲直到完整:在自定义流式处理中,收集所有 chunk,直到生成结束,再一次性解析。这牺牲了流式的即时性,但保证了解析正确。

  3. 增量解析器:对于一些格式(如 JSON Lines),可以实现一个状态机,每收到一行就解析一次。但 LangChain 目前没有内置通用的增量解析器。

  4. 后处理:在不改变 Parser 的情况下,前端可以累积字符,当检测到完整消息(例如 JSON 闭合、特定分隔符)后再渲染。但这样 Parsing 的负担转移到了客户端。

实践建议:对于需要结构化输出的场景,尽量让 LLM 返回简单文本,然后通过后续的链步骤或工具调用进行结构化,而不是依赖流式输出时的 Parser。如果必须流式输出 JSON,推荐使用 JsonOutputToolsParser,因为它基于工具调用的 JSON,天然适合流式。


🧠 在 Agent 中使用流式输出时,你会看到什么?Agent 的 Thought 也会流式吗?

Agent 的流式输出与普通链有所不同,因为 Agent 内部是多步推理的。你会看到每一步的中间结果,而不仅仅是最终答案。

默认行为(使用 AgentExecutor.stream):

  • 每个完整步骤结束后,会产生一个事件,包含该步骤的 ActionFinal 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 方法也可以直接实现打字机效果:

for chunk in chain.stream({"topic": "黑洞"}):
    print(chunk, end="", flush=True)

这种方法不依赖回调,更灵活。


🔍 如果你想在流式输出过程中,基于已生成的内容做一些实时处理(如敏感词过滤),怎么做?

可以利用 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-streamStreamingResponse 自动设置。

  • 客户端连接中断时,生成器会抛出 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 可以接收一个异步生成器,实现服务端推送。

实现步骤:

  1. 确保 LLM 开启 streaming=True,并使用异步客户端(如 AsyncChatOpenAI)。

  2. 定义异步生成器函数:在函数内部使用 async for 遍历 chain.astream(input),将每个 chunk 格式化为 SSE 消息,然后 yield

  3. 路由返回 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 会在线程池中执行,失去了异步的优势,高并发时性能会急剧下降。务必使用 AsyncChatOpenAIChatOpenAI 默认的异步支持。


⏱️ 流式输出中,如果 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 等非生成式组件通常不实现流式,因为:

  1. 语义不符:流式意味着逐步产出数据,而检索是一次性操作:输入查询,返回一组完整文档。没有中间状态。

  2. 性能无益:检索结果通常体量小(几个到几十个文档),一次性返回远比逐文档流式更高效。

  3. 框架设计:LCEL 的 stream 对于非 LLM 组件,通常会退化为一次 invoke 后包装成单元素的迭代器,不会真正流式。

  4. 下游需求:即使 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 把这个结果包装成一个单元素的迭代器。前端会在一瞬间收到全部文本,失去了打字机效果。

变通方案:

  1. 客户端模拟:后端返回完整文本,前端通过 JavaScript 定时逐字渲染,实现假流式。这能改善感知延迟,但首字等待时间不变。

  2. 服务端分块发送:在自定义链中,拿到完整结果后,按句子或固定长度切片,通过生成器逐步 yield,但这并不是真正的流式,且增加了服务端开销。

  3. 换用支持流式的模型:例如从 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),并设置超时,如果消费者在超时内未消费,则主动中止生成,避免浪费。这是流式系统中必须考虑的工程细节。