系列目录:[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执行语法检查,最后综合分析给出答案。同步模式下:

  1. 你发消息后,终端卡死3-5秒没有任何输出;
  2. 突然一次性蹦出几百字的思考+工具调用+结果,像一堵文字墙砸过来;
  3. 如果模型写得很长,你不知道它是在思考还是卡死了。

流式输出就是为了解决这个问题:模型生成一个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]

注意几个关键细节:

  1. idname只在第一个chunk出现,后续chunk不再重复;
  2. arguments是逐字符/逐片段追加的JSON字符串,不是dict,不是完整JSON,中途拿到的片段是无法解析的;
  3. 多工具并行时用index区分:index=0是第一个工具,index=1是第二个,以此类推;
  4. 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,
    ))

核心逻辑总结:

  1. requests.post(stream=True) + resp.iter_lines() 逐行读取SSE流;
  2. 解析data: 前缀,[DONE]终止;
  3. content直接字符串拼接,每收到一段就yield ContentDelta
  4. tool_calls按index维护builder,第一次见到index时创建builder并yield ToolCallStart,arguments片段追加到args_str并yield ToolCallArgsDelta
  5. 流结束后json.loads(args_str)解析出最终参数dict,组装ToolCall列表,yield StreamDone

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.contentStreamDone里的完整文本(给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。
【完成】模型已输出最终回答

区别在哪?

  1. 思考过程逐字实时输出,没有3秒静默等待;
  2. 工具名在模型决定调用的那一刻就显示出来,不用等arguments完整生成;
  3. 不需要事后截断打印,因为文字是流动出来的,不是砸下来的。

同步和流式的最终产出一致

这是设计上最重要的保证:chat_completion返回LLMResponsechat_completion_stream最后yield的StreamDone.response也是LLMResponse,结构完全一致。

这意味着:

  • 你可以在任何需要的地方从流式切回同步,或从同步切到流式;
  • main.py的TOOL_EXECUTING、VALIDATING等后续状态逻辑零改动,因为它们拿到的还是同一个LLMResponse
  • 写测试时可以用同步模式(结果确定,方便断言),CLI运行时用流式模式(体验好)。

这和Go里io.Reader的设计哲学一致:不管数据是一次性来的还是分批来的,最终消费到的是同样的字节流。

小结

本阶段做的事情很纯粹:

  1. llm_client.py新增chat_completion_stream生成器函数,解析SSE流,按事件yield;
  2. 定义5种StreamEvent子类(ContentDelta/ToolCallStart/ToolCallArgsDelta/StreamDone/StreamError);
  3. 流式tool_calls组装:按index维护builder字典,arguments逐字符追加,流结束后json.loads;
  4. main.py的THINKING状态改用流式消费,实时打印思考过程和工具调用意图;
  5. 同步chat_completion完整保留,Judge校验继续用同步。

下一篇链接:【Java/Go后端手撸原生Agent(第七篇):Token预算管理 + 滑动窗口上下文裁剪】

Logo

Agent 垂直技术社区,欢迎活跃、内容共建。

更多推荐