【Java/Go后端手撸原生Agent(第六篇):流式输出——让Agent思考过程“看得见”】
系列目录:[Phase 1 基础框架] →[Phase 2 结构化+死循环问题] →[Phase 3 状态机+文件读取] → [Phase 4 原生Function Calling+Judge校验] → [Phase 5 多工具并行+BashTool] → Phase 6 流式输出(本文)
上一篇链接:Java/Go后端手撸原生Agent(第五篇)
同步模式的痛点
在前7个Phase中,我们的Agent一直用同步模式调用LLM:发请求→阻塞等待→拿到完整响应→处理。这种方式逻辑简单、容易调试,但有一个体验问题:等待期间用户什么都看不到。
想象这个场景:你问Agent"帮我分析一下这个500行的Python文件里有什么问题",Agent决定先调read_file读文件,再调bash执行语法检查,最后综合分析给出答案。同步模式下:
- 你发消息后,终端卡死3-5秒没有任何输出;
- 突然一次性蹦出几百字的思考+工具调用+结果,像一堵文字墙砸过来;
- 如果模型写得很长,你不知道它是在思考还是卡死了。
流式输出就是为了解决这个问题:模型生成一个token,你就看到一个token。就像ChatGPT网页版那样的打字机效果。
第一性原理:流式到底是什么?
本质上是HTTP协议层的能力,不是什么黑科技。
同步调用的HTTP响应是:客户端发请求 → 服务器完整生成 → 一次性返回整个JSON → 连接关闭。
流式调用(SSE,Server-Sent Events)的HTTP响应是:客户端发请求(带stream: true)→ 服务器开始逐chunk返回 → 每个chunk是一个data: {json}\n\n片段 → 最后一个chunk是data: [DONE]\n\n → 连接关闭。
OpenAI的SSE流格式:
data: {"id":"chatcmpl-xxx","choices":[{"delta":{"role":"assistant"},"index":0}]}
data: {"id":"chatcmpl-xxx","choices":[{"delta":{"content":"你"},"index":0}]}
data: {"id":"chatcmpl-xxx","choices":[{"delta":{"content":"好"},"index":0}]}
data: {"id":"chatcmpl-xxx","choices":[{"delta":{"content":","},"index":0}]}
...
data: [DONE]
每个chunk里的delta字段告诉你"这一帧新增了什么"。content增量就是delta.content,工具调用增量就是delta.tool_calls。
流式tool_calls:最容易踩坑的点
content的流式拼接很简单——就是字符串追加。但tool_calls的流式组装远比content复杂,这是这篇文章的核心知识点。
在同步模式下,API返回的tool_calls是完整的:
{
"tool_calls": [{
"id": "call_abc123",
"type": "function",
"function": {
"name": "calculator",
"arguments": "{\"expr\": \"(6+9)*4\"}"
}
}]
}
但在流式模式下,arguments不是一次性给你的,而是逐字符拼接出来的:
chunk1: delta.tool_calls = [{index:0, id:"call_abc", function:{name:"calculator", arguments:""}}]
chunk2: delta.tool_calls = [{index:0, function:{arguments:'{"ex'}}]
chunk3: delta.tool_calls = [{index:0, function:{arguments:'pr":"('}}]
chunk4: delta.tool_calls = [{index:0, function:{arguments:'6+9)'}}]
chunk5: delta.tool_calls = [{index:0, function:{arguments:'}'}}]
...
[DONE]
注意几个关键细节:
id和name只在第一个chunk出现,后续chunk不再重复;arguments是逐字符/逐片段追加的JSON字符串,不是dict,不是完整JSON,中途拿到的片段是无法解析的;- 多工具并行时用
index区分:index=0是第一个工具,index=1是第二个,以此类推; index在每个chunk里都会出现,用来告诉你"这个delta属于哪个工具"。
所以组装策略是:按index维护builder字典。
代码实现
1. 事件模型
流式输出的本质是"一连串事件",我们用继承来区分事件类型:
class StreamEvent:
"""流式事件基类"""
pass
class ContentDelta(StreamEvent):
"""content文本增量"""
def __init__(self, delta: str):
self.delta = delta
class ToolCallStart(StreamEvent):
"""工具调用开始(已知id和name,arguments还没生成完)"""
def __init__(self, index: int, tool_call_id: str, name: str):
self.index = index
self.tool_call_id = tool_call_id
self.name = name
class ToolCallArgsDelta(StreamEvent):
"""arguments JSON字符串增量"""
def __init__(self, index: int, delta: str):
self.index = index
self.delta = delta
class StreamDone(StreamEvent):
"""流结束,携带完整LLMResponse(和同步模式返回值等价)"""
def __init__(self, response: LLMResponse):
self.response = response
class StreamError(StreamEvent):
"""流错误"""
def __init__(self, error: str):
self.error = error
为什么用类继承而不是{"type": "content_delta", "data": "..."}这种dict?因为用isinstance做分发更Pythonic,IDE补全更好,而且类型安全。
2. 流式生成器核心
def chat_completion_stream(
messages: list[dict],
tools: Optional[list[dict]] = None,
json_mode: bool = False,
) -> Generator[StreamEvent, None, None]:
headers = {
"Authorization": f"Bearer {LLMConfig.API_KEY}",
"Content-Type": "application/json",
}
req = _ChatRequest(
model=LLMConfig.MODEL_NAME,
messages=messages,
temperature=0.1,
stream=True, # 关键:告诉API返回SSE流
)
if tools:
req.tools = tools
req.tool_choice = "auto"
elif json_mode:
req.response_format = {"type": "json_object"}
body = req.model_dump(exclude_none=True)
# ---- 增量组装状态 ----
full_content = ""
tool_call_builders: dict[int, dict] = {} # {index: {id, name, args_str}}
try:
with requests.post(
f"{LLMConfig.BASE_URL}/chat/completions",
headers=headers,
json=body,
timeout=60,
stream=True, # 关键:requests以流式方式读取
) as resp:
resp.raise_for_status()
for line in resp.iter_lines(decode_unicode=True):
if not line:
continue
if not line.startswith("data: "):
continue
data_str = line[len("data: "):]
if data_str.strip() == "[DONE]":
break
try:
chunk = json.loads(data_str)
except json.JSONDecodeError:
continue
if not chunk.get("choices"):
continue
choice = chunk["choices"][0]
delta = choice.get("delta", {})
# ---- content增量 ----
delta_content = delta.get("content")
if delta_content:
full_content += delta_content
yield ContentDelta(delta_content)
# ---- tool_calls增量 ----
delta_tool_calls = delta.get("tool_calls")
if delta_tool_calls:
for tc_delta in delta_tool_calls:
idx = tc_delta.get("index", 0)
fn_delta = tc_delta.get("function", {})
if idx not in tool_call_builders:
# 第一次见到这个index:新工具调用开始
builder = {
"id": tc_delta.get("id", ""),
"name": fn_delta.get("name", ""),
"args_str": "",
}
tool_call_builders[idx] = builder
yield ToolCallStart(
index=idx,
tool_call_id=builder["id"],
name=builder["name"],
)
# 追加arguments片段
args_delta = fn_delta.get("arguments")
if args_delta:
tool_call_builders[idx]["args_str"] += args_delta
yield ToolCallArgsDelta(index=idx, delta=args_delta)
except requests.exceptions.RequestException as e:
yield StreamError(f"HTTP请求失败: {e}")
return
except Exception as e:
yield StreamError(f"流式解析失败: {type(e).__name__}: {e}")
return
# ---- 流结束:组装最终的LLMResponse ----
final_tool_calls = None
if tool_call_builders:
final_tool_calls = []
for idx in sorted(tool_call_builders.keys()):
b = tool_call_builders[idx]
args_str = b["args_str"] or "{}"
try:
args = json.loads(args_str)
except json.JSONDecodeError:
args = {}
final_tool_calls.append(ToolCall(
id=b["id"],
name=b["name"],
arguments=args,
))
yield StreamDone(LLMResponse(
content=full_content if full_content else None,
tool_calls=final_tool_calls,
))
核心逻辑总结:
requests.post(stream=True)+resp.iter_lines()逐行读取SSE流;- 解析
data:前缀,[DONE]终止; - content直接字符串拼接,每收到一段就yield
ContentDelta; - tool_calls按index维护builder,第一次见到index时创建builder并yield
ToolCallStart,arguments片段追加到args_str并yieldToolCallArgsDelta; - 流结束后
json.loads(args_str)解析出最终参数dict,组装ToolCall列表,yieldStreamDone。
3. main.py改造:THINKING状态用流式
有了流式接口,主循环怎么接?改动非常小——把原来的一行resp = chat_completion(...)换成事件循环:
if state == AgentTaskState.THINKING:
messages = [{"role": "system", "content": SYSTEM_PROMPT}]
messages.extend(memory.get_messages())
print(f"\n=== 第{loop_count}轮 THINKING ===")
resp = None
for event in chat_completion_stream(messages, tools=tools_schema):
if isinstance(event, ContentDelta):
print(event.delta, end="", flush=True) # 打字机效果
elif isinstance(event, ToolCallStart):
print(f"\n【工具调用意图】{event.name}") # 即时显示工具名
elif isinstance(event, ToolCallArgsDelta):
pass # CLI下不打印参数片段
elif isinstance(event, StreamDone):
resp = event.response # 拿到完整响应
elif isinstance(event, StreamError):
print(f"\n【LLM错误】{event.error}")
state = AgentTaskState.THINKING
break
if resp is None:
continue # 流异常,重试下一轮
print() # 流式输出后换行
# ---- 下面的逻辑和同步模式完全一样 ----
if resp.has_tool_calls:
tool_calls_dicts = [tc.to_openai_dict() for tc in resp.tool_calls]
memory.add_assistant(content=resp.content, tool_calls=tool_calls_dicts)
pending_tool_calls = list(resp.tool_calls)
state = AgentTaskState.TOOL_EXECUTING
continue
else:
pending_draft_answer = resp.content
state = AgentTaskState.VALIDATING
continue
关键设计决策:
同步接口保留。 Judge校验仍然用chat_completion(json_mode=True)——Judge是内部短调用,不需要打字机效果,用同步代码更简洁。不是所有地方都要流式,同步和流式是互补的,不是替代关系。
展示层和数据层分离。 流式打印是给人看的(CLI展示),写入memory的resp.content是StreamDone里的完整文本(给LLM下一轮看的)。不要因为截断了终端打印就截断memory——以前同步模式截断200字符打印也是一样的道理,截断只影响print,不影响存入memory的内容。
流式模式不需要截断思考文本。 同步模式截断是因为"等3-5秒一墙文字砸下来",流式模式文字逐字出来,一旦模型决定调工具,ToolCallStart立即触发,思考自然停止,不存在"一墙文字"问题。
运行效果对比
Phase 7(同步):
=== 第1轮 THINKING ===
(3秒静默...)
【推理思考】用户需要计算(6+9)*4,我需要使用calculator工具...
【工具调用】calculator({"expr": "(6+9)*4"})
【工具结果】计算结果: (6+9)*4 = 60
=== 第2轮 THINKING ===
(2秒静默...)
【最终回答】(6+9)*4 = 60
Phase 8(流式):
=== 第1轮 THINKING ===
用户需要计算(6+9)*4,我先使用calculator工具。
【工具调用意图】calculator
【工具调用】calculator({"expr": "(6+9)*4"})
【工具结果】计算结果: (6+9)*4 = 60
=== 第2轮 THINKING ===
根据计算结果,(6+9)*4 = 60。
【完成】模型已输出最终回答
区别在哪?
- 思考过程逐字实时输出,没有3秒静默等待;
- 工具名在模型决定调用的那一刻就显示出来,不用等arguments完整生成;
- 不需要事后截断打印,因为文字是流动出来的,不是砸下来的。
同步和流式的最终产出一致
这是设计上最重要的保证:chat_completion返回LLMResponse,chat_completion_stream最后yield的StreamDone.response也是LLMResponse,结构完全一致。
这意味着:
- 你可以在任何需要的地方从流式切回同步,或从同步切到流式;
- main.py的TOOL_EXECUTING、VALIDATING等后续状态逻辑零改动,因为它们拿到的还是同一个
LLMResponse; - 写测试时可以用同步模式(结果确定,方便断言),CLI运行时用流式模式(体验好)。
这和Go里io.Reader的设计哲学一致:不管数据是一次性来的还是分批来的,最终消费到的是同样的字节流。
小结
本阶段做的事情很纯粹:
- 在
llm_client.py新增chat_completion_stream生成器函数,解析SSE流,按事件yield; - 定义5种StreamEvent子类(ContentDelta/ToolCallStart/ToolCallArgsDelta/StreamDone/StreamError);
- 流式tool_calls组装:按index维护builder字典,arguments逐字符追加,流结束后json.loads;
- main.py的THINKING状态改用流式消费,实时打印思考过程和工具调用意图;
- 同步
chat_completion完整保留,Judge校验继续用同步。
更多推荐
所有评论(0)