代码仓库ChainReaction

agent.invoke() 慢是个问题,比慢更难受的是沉默。一个带工具的 agent,一次提问背后至少两次模型调用和一次工具执行。DeepSeek 上这两次调用各要一两秒,工具再花点时间,用户盯着空白界面等五秒,然后一段完整回答突然出现。中间那五秒里 agent 其实一直在干活,只是没人告诉外面。

流式要解决的就是这五秒。它把过程中的事件提前暴露出来:模型开始吐第一个 token 了,它决定调哪个工具、参数是什么,工具返回了什么,第二次模型调用又开始了。这些事件本身就是产品体验的一部分,聊天界面里的”正在调用 get_weather”就是这么来的。

这篇讲清楚 LangChain 里流式的两层 API,以及它们各自适合放在哪。底层是 agent.stream() 和它的 stream_mode,你拿到的是按 mode 切分的事件;上层是 agent.stream_events(version="v3"),同一个运行被投影成几条独立的、带类型的流。读完你应该能判断该用哪一层,以及为什么自己的流式代码会重复打印同一段文字。

一次 agent 运行里到底发生了什么

先把时间线摆出来。我用一个只有 get_weather 的 agent,问”北京今天天气怎么样”。DeepSeek 第一次调用返回的是一个 tool_call,正文一个字都没有;agent 执行工具拿到结果,把 ToolMessage 拼进消息列表,再发起第二次模型调用,这次才产出给人看的回答。

invoke() 只给你最后一行。stream() 把上面每一个箭头都变成可以消费的事件,区别只是切法不同。

这件事在工程上的意义比打字机效果大。终端里逐 token 打印只是最表面的用法,真正需要流式的场景是中间状态要被程序消费:前端要显示”正在调用 get_weather”,后端想在模型决定调工具时先做一次权限校验,监控要记录每一步的耗时。这些场景要的是事件,文字只是顺带。

还有一点要提前说清:工具执行本身是普通 Python 函数调用,它的返回值不会逐字符吐出来。流式覆盖的是模型这一侧,包括模型决定调哪个工具、以及模型生成的每一个 token。工具那一段的进度得靠工具自己往外写,后面 custom 模式会讲。

四个 stream_mode 分别吐什么

stream_mode 决定切法。单次只传一个字符串,chunk 就是对应的载荷;传列表就同时开多路,这时 chunk 会变成 (mode, payload) 二元组(旧格式)或带 type 字段的字典(version="v2")。

mode 每次吐出的东西 适合拿它做什么
values 这一步结束后的完整 state,state["messages"] 是截至当前的全部消息 想直接渲染完整对话,或者调试时想看全量状态
updates 只含本步被改动的键,外层带节点名,例如 {"model": {"messages": [...]}} 判断 agent 走到哪一步,拿到已解析的 tool_calls 和 ToolMessage
messages (AIMessageChunk, metadata) 二元组,来自每一次模型调用 逐 token 渲染,metadata 里有 langgraph_node
custom 节点里 get_stream_writer() 写出的任意对象 进度提示、检索进度、业务状态

同一个运行的四种切法长这样:

values:每步一份全量 state

模型实例按系列约定统一用 DeepSeek。这个配置后面所有脚本都复用:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
import os

from langchain.agents import create_agent
from langchain_openai import ChatOpenAI

model = ChatOpenAI(
api_key=os.getenv("DEEPSEEK_API_KEY"),
base_url="https://api.deepseek.com/v1",
model="deepseek-chat",
temperature=0.1,
max_tokens=1000,
)


def get_weather(city: str) -> str:
"""查询某个城市今天的天气。"""
return f"{city}今天晴,22 摄氏度,微风。"


agent = create_agent(model=model, tools=[get_weather])

question = {"messages": [{"role": "user", "content": "北京今天天气怎么样?用一句话回答。"}]}

for step, state in enumerate(agent.stream(question, stream_mode="values")):
messages = state["messages"]
last = messages[-1]
print(f"--- step {step} | 消息总数 {len(messages)} | 最后一条 {type(last).__name__} ---")
print(f" text = {last.text!r}")
if getattr(last, "tool_calls", None):
print(f" tool_calls = {last.tool_calls}")

真实输出(Streaming/StreamValues.py):

1
2
3
4
5
6
7
8
9
--- step 0 | 消息总数 1 | 最后一条 HumanMessage ---
text = '北京今天天气怎么样?用一句话回答。'
--- step 1 | 消息总数 2 | 最后一条 AIMessage ---
text = ''
tool_calls = [{'name': 'get_weather', 'args': {'city': '北京'}, 'id': 'call_00_9LjxS7eTOXBGq3s9DKLy1337', 'type': 'tool_call'}]
--- step 2 | 消息总数 3 | 最后一条 ToolMessage ---
text = '北京今天晴,22 摄氏度,微风。'
--- step 3 | 消息总数 4 | 最后一条 AIMessage ---
text = '北京今天天气晴朗,气温22摄氏度,伴有微风。'

注意 step 0:第一次 yield 时模型还没被调用,state 里只有你传进去的 HumanMessage。所以 values 模式天然适合做”渲染整个对话”,每一步把消息列表重画一遍即可,不需要自己维护历史。

messages:真正的逐 token

stream_mode="messages" 是 token 级流式的来源。单传一个 mode 时,chunk 直接是 (token, metadata) 二元组:

1
2
3
4
5
6
7
question = {"messages": [{"role": "user", "content": "用一句话介绍杭州。"}]}

count = 0
for token, metadata in agent.stream(question, stream_mode="messages"):
print(f"[{metadata.get('langgraph_node')}] chunk_position={token.chunk_position!r} text={token.text!r}")
count += 1
print(f"共收到 {count} 个 chunk")

真实输出(Streaming/StreamTokens.py,中间省略若干行):

1
2
3
4
5
6
7
8
9
10
11
12
13
[model] chunk_position=None text=''
[model] chunk_position=None text='杭州'
[model] chunk_position=None text='是'
[model] chunk_position=None text='“'
[model] chunk_position=None text='人间'
[model] chunk_position=None text='天堂'
[model] chunk_position=None text='”,'
...
[model] chunk_position=None text='厚重'
[model] chunk_position=None text='。'
[model] chunk_position=None text=''
[model] chunk_position='last' text=''
共收到 28 个 chunk

两个细节值得记住。第一个 chunk 的 text 是空字符串,DeepSeek 会先发一个只带 role 的 delta,直接拿它当正文会多出一次空渲染。最后一个 chunk 的 chunk_position 是 "last",这是判断”这次模型调用结束了”的官方信号。

把 chunk 拼回完整消息

metadata 里除了 langgraph_node,还可能出现 lc_agent_name。多 agent 场景下给每个 create_agent 传一个 name,token 的归属就能区分开,否则前端只能把几个 agent 的输出糊在一起。嵌套 agent 作为工具被调用时,父图不会自动透传内层的 token,得在 stream() 上额外传 subgraphs=True。

AIMessageChunk 支持相加,加出来的还是 chunk,usage_metadata 也是在这时候才有值:

1
2
3
4
5
6
7
8
9
10
full = None
for token, metadata in agent.stream(question, stream_mode="messages"):
if metadata.get("langgraph_node") != "model":
continue
full = token if full is None else full + token
if token.chunk_position == "last":
print(f"拼接结果: {full.text}")
print(f"类型: {type(full).__name__}")
print(f"usage: {full.usage_metadata}")
full = None
1
2
3
拼接结果: 杭州是“人间天堂”,既有西湖的柔美诗意,也有龙井的清香与千年宋韵的厚重。
类型: AIMessageChunk
usage: {'input_tokens': 9, 'output_tokens': 25, 'total_tokens': 34, 'input_token_details': {'cache_read': 0}, 'output_token_details': {}}

updates 与 custom

updates 按节点切,一次 agent 运行通常看到三条:model 节点给出带 tool_calls 的 AIMessage,tools 节点给出 ToolMessage,model 节点再给出最终回答。它的真实输出形态在后面工具那一节有完整例子。

这里先把 values 和 updates 的取舍说清楚。values 每一步都把整个 state 重新 yield 一次,对话越长,重复发送的消息越多,走到第 20 步时拿到的 chunk 里装着前 19 步的全部内容。updates 只发本步新增的那条消息,长对话下省得多。代价是消息列表要自己维护,把每个节点的增量拼起来。做聊天界面我一般用 updates 配 messages,写调试脚本用 values 更省事。

custom 模式要配合 get_stream_writer(),稍后单独讲。

顺带一个实战里会撞上的东西:人在回路(human-in-the-loop)的审批请求也从 updates 里出来,键名是 __interrupt__。拦截到的对象带着 action_requests,前端据此弹确认框,用户点完之后用 Command(resume=...) 接着跑同一个 stream 循环。

chunk 的结构:v1 与 v2 两种格式

单 mode 时 chunk 就是载荷本身。多 mode 时格式取决于 version 参数。不传 version 走旧格式,chunk 是 (mode, payload) 元组;传 version="v2" 后每个 chunk 是统一的字典,带 type、ns、data 三个键,不用再判断”这次拿到的是元组还是值”。

ns 是命名空间元组,根图是空元组 (),子图里会带上路径。文档说 v2 需要 LangGraph 1.1 以上,我本机是 1.2.12,实测可用。

多 mode 的实战写法,同时盯 token 和完整消息:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
from langchain.messages import AIMessage, AIMessageChunk, ToolMessage

question = {"messages": [{"role": "user", "content": "上海今天天气怎么样?用一句话回答。"}]}

for chunk in agent.stream(question, stream_mode=["messages", "updates"], version="v2"):
if chunk["type"] == "messages":
token, metadata = chunk["data"]
if not isinstance(token, AIMessageChunk):
continue
if token.tool_call_chunks:
for tc in token.tool_call_chunks:
print(f"[messages/{metadata.get('langgraph_node')}] tool_call_chunk: {tc}")
elif token.text:
print(f"[messages/{metadata.get('langgraph_node')}] text: {token.text!r}")
else:
for source, update in chunk["data"].items():
last = update["messages"][-1]
if isinstance(last, AIMessage) and last.tool_calls:
print(f"[updates/{source}] 完成的 tool_calls: {last.tool_calls}")
elif isinstance(last, ToolMessage):
print(f"[updates/{source}] 工具返回: {last.text}")

这里有个必须写的判断:messages 模式不只会吐 AIMessageChunk,工具执行完之后 ToolMessage 也会以 message 的形式出现,而 ToolMessage 没有 tool_call_chunks 属性。我第一版脚本没加 isinstance 检查,直接抛了 AttributeError: 'ToolMessage' object has no attribute 'tool_call_chunks'。文档里的示例都带着这个判断,抄的时候别省。

工具调用的流式细节

上面那段跑出来是这样(Streaming/StreamToolCalls.py):

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
[messages/model] tool_call_chunk: {'name': 'get_weather', 'args': '', 'id': 'call_00_Ycao0sposNDVRYsSdDll4692', 'index': 0, 'type': 'tool_call_chunk'}
[messages/model] tool_call_chunk: {'name': None, 'args': '{', 'id': None, 'index': 0, 'type': 'tool_call_chunk'}
[messages/model] tool_call_chunk: {'name': None, 'args': '"', 'id': None, 'index': 0, 'type': 'tool_call_chunk'}
[messages/model] tool_call_chunk: {'name': None, 'args': 'city', 'id': None, 'index': 0, 'type': 'tool_call_chunk'}
[messages/model] tool_call_chunk: {'name': None, 'args': '"', 'id': None, 'index': 0, 'type': 'tool_call_chunk'}
[messages/model] tool_call_chunk: {'name': None, 'args': ': ', 'id': None, 'index': 0, 'type': 'tool_call_chunk'}
[messages/model] tool_call_chunk: {'name': None, 'args': '"', 'id': None, 'index': 0, 'type': 'tool_call_chunk'}
[messages/model] tool_call_chunk: {'name': None, 'args': '上海', 'id': None, 'index': 0, 'type': 'tool_call_chunk'}
[messages/model] tool_call_chunk: {'name': None, 'args': '"', 'id': None, 'index': 0, 'type': 'tool_call_chunk'}
[messages/model] tool_call_chunk: {'name': None, 'args': '}', 'id': None, 'index': 0, 'type': 'tool_call_chunk'}
[updates/model] 完成的 tool_calls: [{'name': 'get_weather', 'args': {'city': '上海'}, 'id': 'call_00_Ycao0sposNDVRYsSdDll4692', 'type': 'tool_call'}]
[updates/tools] 工具返回: 上海今天晴,22 摄氏度,微风。
[messages/model] text: '上海'
[messages/model] text: '今天'
[messages/model] text: '晴'
[messages/model] text: ','
[messages/model] text: '气温'
[messages/model] text: ' '
[messages/model] text: '22'
...

工具参数不是一次性给的。第一片带 name 和 id,之后每片只有 args,而且是 {"city": "上海"} 这段 JSON 的字符级碎片。想在前端显示”正在调用 get_weather”,用第一片里的 name 就够了;想拿到能执行的参数字典,得等 updates 里的 tool_calls,那里的 args 已经是解析好的 {'city': '上海'}。

如果不想依赖 updates,也可以自己聚合:把 chunk 一路相加,等 chunk_position == "last" 时读 full.tool_calls,那里已经是解析好的列表。文档把这叫在流循环里聚合,适合那些没有被写进 state 的消息,比如中间件里临时发起的模型调用。

这也解释了一个常见误解:index 是给同一次模型调用里多个并行工具调用排序用的。模型一次决定调两个工具时,会看到 index: 0 和 index: 1 两组碎片交错出现,按 index 分组拼接才不会串。

两层 API 怎么选

到这里底层这套已经够用,stream_mode 加几个 if 能覆盖大部分终端场景。它有个结构性问题:模式一多,消费端就长成一棵判断树,messages 里混着 token 和 ToolMessage,updates 的键又跟节点名绑在一起。

维度 stream + stream_mode stream_events(version="v3")
输出形态 按 mode 切分的 chunk,多 mode 时是元组或字典 每条 projection 一路独立的、带类型的流
消费方式 循环里对 chunk["type"] 做分支 各自 for 循环,或 interleave、asyncio.gather
token 与工具 token 在 messages,完成的调用在 updates,要跨 mode 拼 message.text 与 stream.tool_calls 各管一段
成熟度 稳定 beta,运行时有 LangChainBetaWarning
适合 终端脚本、简单后端、需要 values 等底层 mode 要同时喂多个 UI 组件的应用

我的建议是终端和调试用 stream,因为 values 看全量状态最直观;面向产品的服务端接口用 v3,省掉那棵判断树。两边不互斥,同一个 agent 两套都能开。

stream_events:v3 的事件模型

stream_mode 的麻烦在于消费端要写一堆 if chunk["type"] == ...。v3 换了个思路:同一次运行被投影成几条独立的流,每条流有自己的类型,互不干扰。文档里推荐新项目直接用这套 API,它从 LangChain 1.3 开始提供,当前还是 beta,运行时会打 LangChainBetaWarning: The v3 streaming protocol on Pregel is experimental。

可用的 projection 有 messages、tool_calls、values、subgraphs、extensions,加上最终的 output。单个 message 流还能再往下拆:.text、.reasoning、.tool_calls、.output。.reasoning 在 DeepSeek 上是空的,后面单独讲。

只消费一条流最省事:

1
2
3
4
5
6
7
8
9
stream = agent.stream_events(question, version="v3")

for message in stream.messages:
print(f"[node={message.node}] ", end="")
for delta in message.text:
print(delta, end="", flush=True)
print()

print(f"final: {stream.output['messages'][-1].text}")
1
2
3
[node=model] 
[node=model] 广州今天天气晴朗,气温22摄氏度,伴有微风。
final: 广州今天天气晴朗,气温22摄氏度,伴有微风。

第一行空的是模型第一次调用,它只产出了 tool_call,没有文本,所以 .text 是空的。同一个运行里每次模型调用对应一个 ChatModelStream,不是整个运行合成一条。

工具执行走另一条流:

1
2
3
4
5
6
7
stream = agent.stream_events(question, version="v3")

for call in stream.tool_calls:
deltas = "".join(call.output_deltas)
print(f"[tool_calls] {call.tool_name}({call.input})")
print(f" output_deltas={deltas!r}")
print(f" output={call.output!r} error={call.error!r}")
1
2
3
[tool_calls] get_weather({'city': '广州'})
output_deltas=''
output=ToolMessage(content='广州今天晴,22 摄氏度,微风。', name='get_weather', id='e03a9cf6-f4c1-4915-8cb6-b05fab7fdd73', tool_call_id='call_00_NHl6cnRnnYzZI4c2VMqh1314') error=None

output_deltas 是空的,因为 get_weather 是个瞬间返回的普通函数,没有增量可给。但这个迭代动作不能省:call.output 在 output_deltas 排空之后才被填充。我一开始没排空就直接读 call.output,拿到的是 None。

要同时看几条流,同步代码用 interleave:

1
2
3
4
5
6
7
8
9
10
11
stream = agent.stream_events(question, version="v3")

for name, item in stream.interleave("messages", "tool_calls", "values"):
if name == "messages":
for delta in item.text:
print(f"[messages/{item.node}] {delta}")
elif name == "tool_calls":
"".join(item.output_deltas)
print(f"[tool_calls] {item.tool_name}({item.input}) -> {item.output.text!r}")
else:
print(f"[values] 快照最后一条 {type(item['messages'][-1]).__name__}")
1
2
3
4
5
6
7
8
9
10
11
[values] 快照最后一条 HumanMessage
[values] 快照最后一条 AIMessage
[tool_calls] get_weather({'city': '广州'}) -> '广州今天晴,22 摄氏度,微风。'
[values] 快照最后一条 ToolMessage
[messages/model] 广州
[messages/model] 今天
[messages/model] 天气
[messages/model] 晴朗
...
[values] 快照最后一条 AIMessage
final: 广州今天天气晴朗,气温22摄氏度,伴有微风。

values 的快照顺序和 stream_mode="values" 一致,四条快照对应四步。

一个我踩到的坑

这几条投影是同时从运行里分出来的活流,不是可以事后重放的结果集。如果先把 stream.messages 整个跑完,再去取 stream.tool_calls,会拿到空列表:

1
2
3
4
5
stream = agent.stream_events(question, version="v3")
for _ in stream.messages:
pass
print(f"tool_calls 里剩下的条目: {list(stream.tool_calls)}")
print(f"values 里剩下的条目: {list(stream.values)}")
1
2
tool_calls 里剩下的条目: []
values 里剩下的条目: []

多投影就用 interleave;异步代码用 astream_events 加 asyncio.gather 并发消费。顺序消费两条流,第二条会是空的。

自定义事件:get_stream_writer

custom 模式配合 get_stream_writer() 使用。工具执行到一半想告诉外面进度,就写一条:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
from langgraph.config import get_stream_writer


def lookup_order(order_id: str) -> str:
"""查询订单状态。"""
writer = get_stream_writer()
writer(f"正在查询订单 {order_id}")
writer(f"订单 {order_id} 已发货")
return "已发货,预计明天送达。"


agent = create_agent(model=model, tools=[lookup_order])

question = {"messages": [{"role": "user", "content": "订单 A100 到哪了?用一句话回答。"}]}

for chunk in agent.stream(question, stream_mode=["updates", "custom"], version="v2"):
if chunk["type"] == "custom":
print(f"[custom] {chunk['data']}")
continue
for source, update in chunk["data"].items():
last = update["messages"][-1]
print(f"[updates/{source}] {type(last).__name__} text={last.text!r}")
1
2
3
4
5
[updates/model] AIMessage text=''
[custom] 正在查询订单 A100
[custom] 订单 A100 已发货
[updates/tools] ToolMessage text='已发货,预计明天送达。'
[updates/model] AIMessage text='订单 A100 已发货,预计明天送达。'

writer 接受任意对象,字符串、字典都行。有个限制要提前知道:工具里一旦用了 get_stream_writer,这个工具就只能在 LangGraph 执行上下文里跑,单独 lookup_order("A100") 会失败。写单元测试时要么绕过,要么把 writer 相关逻辑挪到外面。

接到终端和前端

终端消费就是边收边写,flush 不能忘:

1
2
3
4
5
6
7
8
9
import sys

agent = create_agent(model=model, tools=[])
question = {"messages": [{"role": "user", "content": "用一句话介绍苏州。"}]}

for token, metadata in agent.stream(question, stream_mode="messages"):
if token.text:
sys.stdout.write(token.text)
sys.stdout.flush()

真实输出(Streaming/StreamSSE.py):

1
苏州是一座将古典园林的精致与江南水乡的温婉完美融合的千年古城,素有“人间天堂”之美誉。

浏览器不能直接用这个循环,中间要转成 SSE 帧:

1
2
3
4
5
6
7
8
9
10
11
12
import json


def sse(event: dict) -> str:
return "data: " + json.dumps(event, ensure_ascii=False) + "\n\n"


for token, metadata in agent.stream(question, stream_mode="messages"):
if not token.text:
continue
print(sse({"type": "token", "node": metadata.get("langgraph_node"), "text": token.text}))
print(sse({"type": "done"}))
1
2
3
4
5
6
7
8
9
10
data: {"type": "token", "node": "model", "text": "苏州"}

data: {"type": "token", "node": "model", "text": "是一座"}

data: {"type": "token", "node": "model", "text": "将"}

data: {"type": "token", "node": "model", "text": "古典"}

...(省略剩余帧)
data: {"type": "done"}

每帧一个 JSON,前端 EventSource 收到就 append。ensure_ascii=False 别丢,否则中文会变成一堆 \uXXXX。

不想自己写这套协议,可以直接用 LangChain 的前端 SDK。React 里是 useStream,它把 messages、toolCalls、values、interrupts 都做成响应式状态,后端换 createAgent 就行,协议细节由 SDK 处理。

结构化输出和流式的配合

DeepSeek 不支持原生结构化输出,按系列约定走 ToolStrategy。这意味着一件事:结构化输出在底层就是一次工具调用,所以流式看到的全是 tool_call_chunk,校验后的对象要等运行结束。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
from pydantic import BaseModel

from langchain.agents.structured_output import ToolStrategy


class CityInfo(BaseModel):
"""城市信息。"""

city: str
famous_for: str


agent = create_agent(model=model, tools=[], response_format=ToolStrategy(CityInfo))

question = {"messages": [{"role": "user", "content": "介绍杭州,给出城市名和它最有名的一样东西。"}]}

for token, metadata in agent.stream(question, stream_mode="messages"):
for block in token.content_blocks:
print(block)
1
2
3
4
5
6
7
8
9
{'type': 'tool_call_chunk', 'id': 'call_00_YKrwbh5TtVv5ERlnfobb8475', 'name': 'CityInfo', 'args': '', 'index': 0}
{'type': 'tool_call_chunk', 'id': None, 'name': None, 'args': '{', 'index': 0}
{'type': 'tool_call_chunk', 'id': None, 'name': None, 'args': '"', 'index': 0}
{'type': 'tool_call_chunk', 'id': None, 'name': None, 'args': 'city', 'index': 0}
...
{'type': 'tool_call_chunk', 'id': None, 'name': None, 'args': '西湖', 'index': 0}
{'type': 'tool_call_chunk', 'id': None, 'name': None, 'args': '"', 'index': 0}
{'type': 'tool_call_chunk', 'id': None, 'name': None, 'args': '}', 'index': 0}
{'type': 'text', 'text': "Returning structured response: city='杭州' famous_for='西湖'"}

流完这些碎片之后,才在最终 state 的 structured_response 里拿到 Pydantic 对象:

1
2
city='杭州' famous_for='西湖'
CityInfo

所以前端能做的只有”正在生成结构化结果”这类提示,没法逐字段渲染一个还没校验的对象。

同样的道理适用于工具返回值。get_weather 是普通函数,它返回的字符串在 ToolMessage 里一次性出现,没有 token 级过程可流。想让工具输出也流起来,得让工具自己调 writer,或者让工具内部再调一次模型。

还有一些东西压根不在流里。disable_streaming=True(部分集成上写作 streaming=False)的模型不产 token,整个运行只在最后给你一条完整消息,多 agent 系统里常用它控制哪个 agent 该往前端吐字。checkpointer 写盘的最终 state 也不属于事件流,它只在运行结束后可查。

结果就是前端能做的判断很有限:工具执行期间只能显示”调用中”,结构化输出只能显示”生成中”。这不是缺陷,是数据本身的形状决定的。

reasoning:能数出长度,读不到正文

前面所有例子都用 deepseek-chat。它没有推理输出,原始 SSE 里连 reasoning_content 这个字段都没有,所以 reasoning 一直是”照文档写的,没验证过”。换 deepseek-flash 重跑,情况清楚了。

先看模型这一侧。绕开 LangChain,直接用 openai 客户端读原始 SSE:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
import os

from openai import OpenAI

client = OpenAI(
api_key=os.getenv("DEEPSEEK_API_KEY"),
base_url="https://api.deepseek.com/v1",
)

stream = client.chat.completions.create(
model="deepseek-flash",
messages=[{"role": "user", "content": "9.11 和 9.9 哪个大?说清理由。"}],
stream=True,
)

n_reasoning, n_content = 0, 0
for chunk in stream:
if not chunk.choices:
continue
delta = chunk.choices[0].delta.model_dump(exclude_none=True)
if delta.get("reasoning_content"):
if n_reasoning < 3:
print(f"[reasoning] {delta['reasoning_content']!r}")
n_reasoning += 1
if delta.get("content"):
if n_content < 5:
print(f"[content] {delta['content']!r}")
n_content += 1

print(f"reasoning 分片 {n_reasoning} 个,content 分片 {n_content} 个")

真实输出(Streaming/ReasoningRawSSE.py):

1
2
3
4
5
6
7
8
9
[reasoning] '我们需要'
[reasoning] '回答'
[reasoning] '中文'
[content] '如果'
[content] '按'
[content] '**'
[content] '普通'
[content] '十进制'
reasoning 分片 364 个,content 分片 197 个

每次采样的长度都不一样,这里要看的是字段在不在、正文长什么样。reasoning_content 是一段完整的自然语言推理,先把 364 片流完,content 才开始。

问题出在 LangChain 这一侧。同一个 deepseek-flash 交给 ChatOpenAI,stream_mode="messages" 跑一遍(Streaming/ReasoningStream.py):

1
2
3
4
5
6
7
8
9
收到的 AIMessageChunk 总数: 912
text 非空的分片: 248 text 为空的分片: 664
第一个非空 text 出现在第 663 个分片,它之前全是空分片
additional_kwargs 非空的分片: 0
所有 chunk 出现过的 content_blocks 类型: ['text']
所有 chunk 出现过的 additional_kwargs 键: []
聚合后 usage_metadata: {'input_tokens': 47, 'output_tokens': 910, 'total_tokens': 957, 'input_token_details': {'cache_read': 0}, 'output_token_details': {'reasoning': 661}}
聚合后 content_blocks 类型: ['text']
聚合后 .text 长度: 383

逐条对一下前面那几个问题:

  • content_blocks 里没有 type == "reasoning" 的块。912 个分片出现过的块类型只有 text 一种。
  • additional_kwargs 全是空的,reasoning_content 连键名都没留下。
  • usage_metadata.output_token_details 里拿到了 {'reasoning': 661},reasoning token 数是可观测的。
  • 664 个空分片不是白来的。推理阶段每个 token 都照样走一遍流,只是转换时只读 content,reasoning_content 被扔掉,落下来就是一个 text 为空的 chunk。661 个 reasoning token 对上正文前的 662 个空分片,多出来的那个是模型开头只带 role 的 delta。数量几乎严丝合缝。

stream_events(version="v3") 这边一样。ChatModelStream 确实有 .reasoning 这个投影,但它是空的:

1
2
[node=model] 有 .reasoning 属性: True
text 长度=361 reasoning 分片数=0 reasoning 长度=0

原因写在 langchain-openai 的源码注释里。ChatOpenAI 只针对 OpenAI 官方规范,第三方 provider 加的非标准字段(点名了 reasoning_content、reasoning_details)不会被提取或保留,要完整支持得用 provider 专用包,比如 ChatDeepSeek、ChatOpenRouter。langchain-core 那边的接收端早就准备好了,_extract_reasoning_from_additional_kwargs 会在 additional_kwargs["reasoning_content"] 存在时合成一个 {"type": "reasoning", "reasoning": ...} 块,函数注释里列的就是 Ollama、DeepSeek、XAI、Groq。缺的是中间那一环,provider 包没往 additional_kwargs 里填。

补上这一环会怎样,造个假模型就能看。它 _stream 出来的 chunk 只比平时多带一个 additional_kwargs={"reasoning_content": ...}:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
from langchain_core.language_models import BaseChatModel
from langchain_core.messages import AIMessage, AIMessageChunk
from langchain_core.outputs import ChatGeneration, ChatGenerationChunk, ChatResult


class StubReasoningModel(BaseChatModel):
"""只用来验证通路:像 ChatDeepSeek 那样把 reasoning 写进 additional_kwargs。"""

@property
def _llm_type(self) -> str:
return "stub-reasoning"

def _generate(self, messages, stop=None, run_manager=None, **kwargs):
return ChatResult(
generations=[
ChatGeneration(
message=AIMessage(
content="9.9 更大。",
additional_kwargs={"reasoning_content": "先比整数部分,再比十分位。"},
)
)
]
)

def _stream(self, messages, stop=None, run_manager=None, **kwargs):
pieces = [("9.9", "先比整数部分"), (" 更大。", ",再比十分位。")]
for text, reasoning in pieces:
yield ChatGenerationChunk(
message=AIMessageChunk(
content=text,
additional_kwargs={"reasoning_content": reasoning},
)
)

同一个 agent 换上它,两层 API 都拿到了 reasoning:

1
2
3
4
5
6
7
8
9
-- stream_mode="messages" --
chunk content_blocks: [{"type": "reasoning", "reasoning": "先比整数部分"}, {"type": "text", "text": "9.9"}]
chunk content_blocks: [{"type": "reasoning", "reasoning": ",再比十分位。"}, {"type": "text", "text": " 更大。"}]
聚合后 content_blocks: [{"type": "reasoning", "reasoning": "先比整数部分,再比十分位。"}, {"type": "text", "text": "9.9 更大。"}]
聚合后 .text: '9.9 更大。'
-- stream_events(version="v3") --
[node=model] reasoning 分片数=2
reasoning 正文: '先比整数部分,再比十分位。'
text: '9.9 更大。'

reasoning 分片和 text 分片并排躺在同一个 chunk 的 content_blocks 里,靠 type 区分。message.text 只拼 text 块,reasoning 一个字都不掺。v3 下它落在 stream.messages 这条流的 .reasoning 上,.text 和 .reasoning 各自迭代互不干扰。chunk 相加也只管把同名块接起来,推理正文和最终答案不会串到一起。

最后拿 deepseek-chat 做个对照,它的原始 SSE 里一片 reasoning_content 都没有:

1
2
3
4
delta 里带 reasoning_content 的分片数: 0
delta 里带 content 的分片数: 148
usage.completion_tokens_details.reasoning_tokens: None
聚合后 usage_metadata: {'input_tokens': 21, 'output_tokens': 131, 'total_tokens': 152, 'input_token_details': {'cache_read': 0}, 'output_token_details': {}}

output_token_details 是空的 {},连计数都没有。所以结论分两档:deepseek-flash 在 LangChain 里能观测到 reasoning token 数,拿不到 reasoning 正文;deepseek-chat 两样都没有。前端能显示”正在思考”,想显示思考内容,得换 provider 包,或者绕开 LangChain 直接读原始 SSE。

常见坑

同一段文字打印两遍是最高频的问题。同时开 messages 和 updates 时,最终回答会以两种形态各来一次:token 逐片来一次,完整的 AIMessage 在 updates 里再来一次。上面 custom 那段输出里的 [updates/model] AIMessage text='订单 A100 已发货,预计明天送达。' 就是完整文本。渲染时按节点和消息类型分工,别两边都往同一个列表里塞。

空 chunk 容易被当成内容。DeepSeek 首个 delta 的 text 是空串,末尾还有一个 chunk_position='last' 的空 chunk;换成 deepseek-flash 之后,整个推理阶段的每个 token 都落成一片空 chunk,一次回答里空片能占到七成。判断有没有正文要用 if token.text,别用 is not None。

content_blocks 上不要做假设。一个 chunk 的 content_blocks 可能是 [],也可能同时带 text 和 tool_call_chunk。按 block["type"] 分派,别取第一个块就当文本处理。

isinstance 过滤不能省。messages 模式里混着 AIMessageChunk 和 ToolMessage,属性不完全一致,前面那个 AttributeError 就是这么来的。

version 和元组解包要对上。传了 version="v2" 之后 chunk 是字典,继续写 for mode, payload in ... 会直接报错;多 mode 又不传 version 时,拿到的确实是元组。这一点容易记混。

小结

  • 一次带工具的 agent 运行至少两次模型调用,invoke 把中间过程全藏起来,流式把它们暴露出来。
  • stream_mode 是切法:values 给全量快照,updates 给节点级增量,messages 给 token,custom 给自己写的事件。
  • token 级流式用 stream_mode="messages",chunk 是 AIMessageChunk,可相加,chunk_position == "last" 表示这次模型调用结束。
  • stream_events(version="v3") 把运行投影成 messages、tool_calls、values 等独立流,用 interleave 或 asyncio.gather 同时消费,顺序消费第二条流会拿到空结果。该 API 目前是 beta。
  • 流不到的东西要认:工具返回值、结构化输出的校验对象、ChatOpenAI 下的 reasoning 正文。deepseek-flash 只在 usage_metadata 里给出 reasoning token 数,正文要换 provider 包,或者绕开 LangChain 读原始 SSE。前端能做的只是”进行中”提示。