七夜zippoe头像
关注
Agent 流式输出实战:SSE/WebSocket 实现实时响应与中间状态展示封面图

Agent 流式输出实战:SSE/WebSocket 实现实时响应与中间状态展示

摘要: 在 AI Agent 应用中,传统的请求-响应模式已无法满足用户对实时反馈的期待。本文系统讲解如何使用 SSE(Server-Sent Events)和 WebSocket 两种技术方案实现 Agent 的流式输出,涵盖后端架构设计、前端打字机效果渲染、Markdown 实时解析、中断处理策略,以及一个完整的前后端实战案例。文章深入对比了 SSE 与 WebSocket 的适用场景,提供了生产级的代码实现,并给出了明确的风险边界。无论你正在构建对话式 Agent 还是复杂的多步骤 Agent 工作流,本文都能为你提供可直接落地的技术方案。

版本声明: 本文基于 Python 3.11+、FastAPI 0.110+、Python-SocketIO 5.11+、前端原生 JavaScript / React 18 编写。相关方案同样适用于 Node.js / Go 等后端语言,核心思路一致。文章中涉及的 LLM API 以 OpenAI 兼容接口为例,兼容 DeepSeek、通义千问等主流模型。阅读本文前,建议你已具备基本的 Web 开发经验和 AI Agent 构建基础。

文章目录


一、为什么 Agent 需要流式输出:用户体验与交互设计

1.1 传统请求-响应模式的痛点

想象一个场景:用户向 Agent 发送一条指令——“帮我分析这份财报,总结核心数据并给出投资建议”。在传统的同步模式下,用户面对的是一个漫长的等待过程:请求发出后,页面转圈,几秒甚至十几秒后才返回完整结果。这段等待期间,用户无法知道 Agent 到底在做什么、进展到哪一步、是否遇到了错误。

这种体验有几个严重的问题:

  • 感知延迟高:LLM 生成一个长回答可能需要 10-30 秒,用户在等待期间完全看不到任何反馈
  • 不确定性焦虑:用户不知道 Agent 是否正常工作,是否会超时失败
  • 无法中断:即使用户发现方向不对,也无法中途叫停,只能等整个响应完成后才能重新提问
  • 缺乏过程透明度:Agent 的中间思考步骤(如搜索、工具调用、推理链)对用户不可见

1.2 流式输出带来的交互变革

流式输出(Streaming Output)将 Agent 的响应过程从"一次性交付"变为"持续交付"。它的核心价值体现在三个层面:

用户发送指令

Agent开始处理

流式输出中间状态

用户实时看到进展

用户可以中断/调整

Agent完成响应

用户获得完整结果

第一层:Token 级流式——LLM 每生成一个 Token 就推送给前端,用户看到文字逐字出现,就像有人在实时打字。这大幅降低了感知延迟,从"等 10 秒看一大段文字"变成"立刻看到第一个字"。

第二层:步骤级流式——当 Agent 执行多步骤任务(如"搜索→分析→总结"),每个步骤的开始、进展、完成都实时推送。用户能看到 Agent 正在执行什么操作,而不是面对一个黑盒。

第三层:交互级流式——在流式过程中,用户可以发送新的指令来影响 Agent 的后续行为。比如用户看到前半段输出方向不对,可以中途发送"换个角度分析",Agent 能感知并调整。

1.3 流式输出 vs 传统模式的对比

维度传统请求-响应流式输出
首字节延迟等待完整生成(5-30s)接近实时(<1s)
中间状态可见不可见完全透明
中断能力不支持支持,可随时停止
用户体验焦虑等待掌控感强,参与感强
实现复杂度简单中等到复杂
网络协议HTTP 请求-响应SSE 或 WebSocket
服务器资源短连接,快速释放长连接,需管理

在这里插入图片描述

图:流式输出三层架构:前端(EventSource/Fetch接收+打字机渲染)→后端(FastAPI StreamingResponse+异步生成器)→LLM API(流式Token输出)

1.4 什么时候需要流式输出

不是所有场景都需要流式。以下是需要流式的典型场景:

  1. 对话式 Agent——用户发送消息,等待 Agent 回复。流式输出让用户更快看到回答。
  2. 任务执行型 Agent——Agent 执行多步骤任务(如写报告、做分析),用户需要看到每一步的进展。
  3. 工具调用型 Agent——Agent 调用外部工具(如搜索、代码执行),用户需要知道调用了什么工具、结果如何。
  4. 协作型 Agent——多个 Agent 协作完成任务,用户需要看到每个 Agent 的工作状态。

1.5 流式输出的技术演进路径

在 AI Agent 的发展历程中,流式输出的实现方式也经历了一个演进过程:

第一阶段:轮询(Polling)

早期的 Agent 应用采用轮询方式——前端每隔 1-2 秒向后端发送一次请求,查询当前进度。这种方式实现简单,但存在明显缺陷:大量无效请求浪费带宽、实时性差(最多 2 秒延迟)、服务器日志被轮询请求淹没。

第二阶段:长轮询(Long Polling)

长轮询是对普通轮询的改进——前端发送请求后,服务器保持连接不立即返回,直到有新数据或超时才返回。前端收到响应后立即发送下一个请求。Comet 技术就是典型的长轮询方案。这种方式减少了无效请求,但每次返回后需要重新建立连接,且 HTTP 头部开销大。

第三阶段:SSE / WebSocket

也就是本文要深入讲解的两种方案。SSE 和 WebSocket 都基于持久连接,一次连接建立后可以持续推送数据,真正实现了实时通信。它们是当前 AI Agent 流式输出的主流选择,也是本文的核心内容。

2015-2018 轮询时代\n前端每2秒请求一次\n实时性差,浪费带宽 2018-2021 长轮询时代\nComet技术\n减少无效请求\n但连接开销大 2021-2023 SSE/WebSocket时代\n持久连接\n实时推送\n成为主流 2023-Now Agent流式标准\nOpenAI流式API\n中间状态推送\n多步骤流式 Agent 流式输出技术演进

第四阶段:Agent 原生流式协议

随着 Agent 框架的发展,一些更高级的流式协议开始出现。例如 LangChain 的 astream_events API、LangGraph 的事件流、OpenAI Assistants API 的流式事件等。这些协议不仅流式输出 LLM token,还能流式推送工具调用、思考过程、状态变更等结构化事件,让 Agent 的整个执行过程完全透明化。但无论上层的协议如何设计,底层的传输机制仍然是 SSE 或 WebSocket。

在这里插入图片描述


二、SSE 方案:Server-Sent Events 实现流式文本输出

2.1 SSE 原理与特点

SSE(Server-Sent Events)是 HTML5 标准的一部分,它允许服务器通过 HTTP 长连接向客户端推送数据。与 WebSocket 不同,SSE 是单向的——只能从服务器到客户端,但这恰恰是 Agent 流式输出最需要的模式。

SSE 的核心优势在于简洁:它基于标准 HTTP 协议,不需要额外的协议升级,不需要特殊的代理配置,天然支持断线重连。对于 Agent 的流式输出场景(服务器→客户端的单向数据流),SSE 通常是首选方案。

LLM API 后端服务器 前端客户端 LLM API 后端服务器 前端客户端 前端实时渲染文字 流结束,关闭连接 POST /chat/stream (HTTP 请求) 发送 Prompt Token 1 data: {"token": "你"} Token 2 data: {"token": "好"} Token 3 data: {"token": ","} [DONE] data: [DONE]

2.2 SSE 数据格式

SSE 使用纯文本格式,每条消息以 data: 开头,以两个换行符 \n\n 结束:

data: {"type": "token", "content": "你"}\n\n
data: {"type": "token", "content": "好"}\n\n
data: {"type": "step", "name": "搜索完成", "status": "done"}\n\n
data: [DONE]\n\n

2.3 后端实现:FastAPI + SSE

以下是使用 FastAPI 实现 SSE 流式输出的完整后端代码:

import json
import asyncio
from fastapi import FastAPI, Request
from fastapi.responses import StreamingResponse
from openai import AsyncOpenAI

app = FastAPI(title="Agent Streaming API")

# 初始化 LLM 客户端(兼容 OpenAI 接口的模型均可)
client = AsyncOpenAI(
    api_key="your-api-key",
    base_url="https://api.deepseek.com/v1"  # 以 DeepSeek 为例
)


async def agent_stream_generator(user_message: str, request: Request):
    """
    Agent 流式输出生成器
    产出三种类型的事件:
    1. step - 步骤状态更新(如"正在思考"、"正在搜索")
    2. token - LLM 生成的文本片段
    3. done - 流结束标记
    """
    # 第一步:推送"正在思考"状态
    yield f"data: {json.dumps({'type': 'step', 'step': 'thinking', 'status': 'started', 'message': 'Agent 正在思考...'})}\n\n"
    
    await asyncio.sleep(0.3)  # 模拟思考延迟
    
    # 第二步:调用 LLM 并流式传输
    yield f"data: {json.dumps({'type': 'step', 'step': 'thinking', 'status': 'completed', 'message': '思考完成,开始生成回复'})}\n\n"
    
    yield f"data: {json.dumps({'type': 'step', 'step': 'generating', 'status': 'started', 'message': '正在生成回复...'})}\n\n"
    
    try:
        stream = await client.chat.completions.create(
            model="deepseek-chat",
            messages=[
                {"role": "system", "content": "你是一个专业的AI助手,请用中文回答。"},
                {"role": "user", "content": user_message},
            ],
            stream=True,  # 关键:启用流式输出
        )
        
        async for chunk in stream:
            # 检查客户端是否断开连接
            if await request.is_disconnected():
                print("客户端已断开,停止生成")
                break
            
            if chunk.choices[0].delta.content is not None:
                token_data = {
                    "type": "token",
                    "content": chunk.choices[0].delta.content
                }
                yield f"data: {json.dumps(token_data, ensure_ascii=False)}\n\n"
        
        # 流结束
        yield f"data: {json.dumps({'type': 'step', 'step': 'generating', 'status': 'completed'})}\n\n"
        yield "data: [DONE]\n\n"
        
    except Exception as e:
        error_data = {
            "type": "error",
            "message": f"生成过程中出错: {str(e)}"
        }
        yield f"data: {json.dumps(error_data, ensure_ascii=False)}\n\n"


@app.post("/api/chat/stream")
async def chat_stream(request: Request):
    """
    SSE 流式对话接口
    接收用户消息,返回 SSE 流
    """
    body = await request.json()
    user_message = body.get("message", "")
    
    return StreamingResponse(
        agent_stream_generator(user_message, request),
        media_type="text/event-stream",
        headers={
            "Cache-Control": "no-cache",
            "Connection": "keep-alive",
            "X-Accel-Buffering": "no",  # 关闭 Nginx 缓冲,确保实时推送
        }
    )

代码解析: 上面的代码实现了完整的 SSE 流式输出后端。核心要点有三个:第一,agent_stream_generator 是一个异步生成器,通过 yield 逐步产出 SSE 格式的数据,它推送的不只是 LLM 的 token,还有步骤状态变更事件(step 类型),让前端能展示 Agent 当前在做什么。第二,StreamingResponse 是 FastAPI 提供的流式响应类,它将异步生成器的输出以 text/event-stream 的 MIME 类型返回给客户端,这是 SSE 的标准内容类型。第三,X-Accel-Buffering: no 这个响应头至关重要——Nginx 默认会缓冲后端的响应内容,这会导致 SSE 数据不能实时推送到客户端,设置此头可以关闭缓冲。

2.4 前端实现:EventSource 接收 SSE

前端使用浏览器原生的 EventSource API 或 fetch 来接收 SSE 数据。由于 EventSource 只支持 GET 请求,而 Agent 交互通常需要 POST(发送消息体),因此推荐使用 fetch + ReadableStream 方案:

/**
 * SSE 流式对话客户端
 * 使用 fetch + ReadableStream 实现,支持 POST 请求
 */
class AgentSSEClient {
  constructor(url) {
    this.url = url;
    this.controller = null;  // AbortController,用于中断
    this.onToken = null;      // token 回调
    this.onStep = null;       // 步骤状态回调
    this.onError = null;      // 错误回调
    this.onDone = null;       // 完成回调
  }

  async sendMessage(message) {
    this.controller = new AbortController();
    
    try {
      const response = await fetch(this.url, {
        method: "POST",
        headers: { "Content-Type": "application/json" },
        body: JSON.stringify({ message }),
        signal: this.controller.signal,
      });

      if (!response.ok) {
        throw new Error(`HTTP error! status: ${response.status}`);
      }

      const reader = response.body.getReader();
      const decoder = new TextDecoder();
      let buffer = "";

      while (true) {
        const { done, value } = await reader.read();
        if (done) break;

        buffer += decoder.decode(value, { stream: true });
        
        // SSE 数据以 \n\n 分隔
        const lines = buffer.split("\n\n");
        buffer = lines.pop();  // 保留最后不完整的数据

        for (const line of lines) {
          if (!line.startsWith("data: ")) continue;
          
          const data = line.slice(6).trim();
          if (data === "[DONE]") {
            this.onDone?.();
            return;
          }

          try {
            const event = JSON.parse(data);
            switch (event.type) {
              case "token":
                this.onToken?.(event.content);
                break;
              case "step":
                this.onStep?.(event);
                break;
              case "error":
                this.onError?.(event.message);
                break;
            }
          } catch (e) {
            console.warn("Failed to parse SSE data:", data);
          }
        }
      }
    } catch (err) {
      if (err.name === "AbortError") {
        console.log("请求已被用户中断");
      } else {
        this.onError?.(err.message);
      }
    }
  }

  // 中断当前流式请求
  stop() {
    if (this.controller) {
      this.controller.abort();
      this.controller = null;
    }
  }
}

// 使用示例
const client = new AgentSSEClient("/api/chat/stream");
let fullText = "";

client.onToken = (token) => {
  fullText += token;
  document.getElementById("output").textContent = fullText;
  // 滚动到底部
  document.getElementById("output").scrollTop = 
    document.getElementById("output").scrollHeight;
};

client.onStep = (event) => {
  console.log(`步骤: ${event.step} -> ${event.status}`, event.message);
  // 更新 UI 显示当前步骤状态
  const statusEl = document.getElementById("status");
  statusEl.textContent = event.message || `${event.step}: ${event.status}`;
};

client.onError = (msg) => {
  console.error("Agent 错误:", msg);
};

client.onDone = () => {
  console.log("流式输出完成");
};

// 发送消息
client.sendMessage("请解释什么是 Transformer 架构");

代码解析: 这段前端代码封装了一个 AgentSSEClient 类,核心使用 fetch API 配合 ReadableStream 来消费 SSE 流。选择 fetch 而非原生 EventSource 的原因是,EventSource 只支持 GET 请求,无法发送 JSON 消息体。代码中有几个关键设计:AbortController 提供了中断能力,用户可以随时停止生成;buffer 变量处理 TCP 分包问题——一次 read() 可能返回不完整的 SSE 数据,需要缓存拼接后按 \n\n 分隔;事件通过回调函数分发,不同类型的消息(token、step、error)分别处理,使得 UI 更新逻辑清晰。

2.5 SSE 的多步骤 Agent 流式

当 Agent 需要执行多个步骤(如搜索→分析→生成)时,SSE 可以推送完整的执行链路:

async def multi_step_agent_stream(user_message: str, request: Request):
    """
    多步骤 Agent 的流式输出
    每个步骤都有状态推送 + 内容输出
    """
    steps = [
        {"id": "understand", "name": "理解用户意图", "executor": understand_intent},
        {"id": "search", "name": "搜索相关信息", "executor": search_web},
        {"id": "analyze", "name": "分析搜索结果", "executor": analyze_results},
        {"id": "generate", "name": "生成最终回复", "executor": generate_response},
    ]
    
    context = {"user_message": user_message}
    
    for step in steps:
        # 推送步骤开始
        yield f"data: {json.dumps({'type': 'step_start', 'step_id': step['id'], 'step_name': step['name']}, ensure_ascii=False)}\n\n"
        
        # 执行步骤,可能产出中间结果
        async for chunk in step["executor"](context):
            if await request.is_disconnected():
                return
            yield f"data: {json.dumps({'type': 'step_data', 'step_id': step['id'], 'content': chunk}, ensure_ascii=False)}\n\n"
            await asyncio.sleep(0.01)  # 让出事件循环
        
        # 推送步骤完成
        yield f"data: {json.dumps({'type': 'step_end', 'step_id': step['id']}, ensure_ascii=False)}\n\n"
    
    yield "data: [DONE]\n\n"

代码解析: 这段代码展示了一个多步骤 Agent 的流式输出框架。它将 Agent 的执行流程抽象为有序的步骤列表,每个步骤包含三个阶段——开始、执行、结束,对应三种 SSE 事件类型。step_data 事件携带步骤执行过程中产生的中间结果(如搜索到的网页摘要、分析得出的关键点),让用户能实时看到每一步的产出。context 字典在步骤间传递数据,understand_intent 的输出成为 search_web 的输入,依此类推。这种设计让复杂的 Agent 工作流变得透明可观测。

2.6 SSE 连接管理与断线重连

在生产环境中,网络不稳定是常态。SSE 的一个重要优势是原生支持断线重连——浏览器 EventSource API 自动在连接断开后重试。但使用 fetch 实现 SSE 时,需要手动实现重连逻辑。

async def sse_with_reconnect(request: Request, last_event_id: str = None):
    """
    支持断线重连的 SSE 接口
    通过 Last-Event-ID 头部恢复中断的流
    """
    async def event_stream():
        # 检查是否有断点恢复信息
        if last_event_id:
            yield f"data: {json.dumps({'type': 'resuming', 'from_event': last_event_id})}\n\n"
            # 从存储中恢复之前的事件位置
            # 例如从 Redis 中获取已发送的事件 ID 列表
            resume_from = int(last_event_id)
        else:
            resume_from = 0
            yield f"data: {json.dumps({'type': 'start', 'message': '新会话开始'})}\n\n"
        
        event_id = resume_from
        
        while True:
            if await request.is_disconnected():
                break
            
            # 获取下一条事件(从消息队列、数据库等)
            event = await get_next_event(event_id)
            if event is None:
                # 没有新事件,发送心跳保持连接
                yield f": heartbeat\n\n"
                await asyncio.sleep(15)  # 每15秒发送一次心跳
                continue
            
            event_id += 1
            yield f"id: {event_id}\n"
            yield f"data: {json.dumps(event, ensure_ascii=False)}\n\n"
    
    return StreamingResponse(
        event_stream(),
        media_type="text/event-stream",
        headers={
            "Cache-Control": "no-cache",
            "Connection": "keep-alive",
            "X-Accel-Buffering": "no",
        }
    )

代码解析: 这段代码展示了 SSE 断线重连的生产级实现。Last-Event-ID 是 SSE 协议的标准字段——浏览器 EventSource 在重连时会自动在请求头中携带 Last-Event-ID,服务端可以根据它恢复中断的流。代码中 id: {event_id} 行为每个事件分配序号,客户端会记录最后收到的事件 ID。当连接断开后重连,服务端从 resume_from 继续推送。心跳机制(: heartbeat)是 SSE 的注释行——以冒号开头的数据不会被客户端当作消息处理,但它能保持 TCP 连接活跃,防止中间代理超时断开。

2.7 SSE 在 Nginx 反向代理下的配置

这是生产环境中最常见的 SSE 踩坑点。Nginx 默认会缓冲后端的响应,导致 SSE 数据积压后才一次性发送,完全失去实时性:

# Nginx SSE 配置
server {
    listen 80;
    server_name api.youragent.com;

    # SSE 接口专用配置
    location /api/sse/ {
        proxy_pass http://backend:8000;
        
        # 关闭缓冲,关键配置
        proxy_buffering off;
        proxy_cache off;
        
        # 关闭 TCP 缓冲,确保小数据包立即发送
        proxy_set_header Connection "";
        proxy_http_version 1.1;
        proxy_set_header X-Accel-Buffering no;
        
        # SSE 长连接超时设置(建议 5-10 分钟)
        proxy_read_timeout 300s;
        proxy_send_timeout 300s;
        
        # 传递客户端真实 IP
        proxy_set_header X-Real-IP $remote_addr;
        proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for;
    }

    # WebSocket 接口配置
    location /ws/ {
        proxy_pass http://backend:8000;
        proxy_http_version 1.1;
        proxy_set_header Upgrade $http_upgrade;
        proxy_set_header Connection "upgrade";
        proxy_read_timeout 300s;
        proxy_send_timeout 300s;
    }
}

代码解析: Nginx SSE 配置的核心是 proxy_buffering off——关闭响应缓冲后,后端产生的每个 SSE 事件会立即转发给客户端。proxy_read_timeout 300s 将读取超时设为 5 分钟,防止 Nginx 在 Agent 思考期间断开连接(默认超时只有 60 秒)。WebSocket 配置中的 Upgrade 和 Connection 头部是协议升级所必需的,缺少任何一个都会导致 WebSocket 握手失败。
在这里插入图片描述

图:左侧SSE单向通信(服务器→客户端数据流,客户端需新建连接发消息),右侧WebSocket双向通信(实时双向箭头,随时发送和接收)

2.8 SSE 的事件类型设计

一个设计良好的 Agent SSE 事件系统应该包含哪些事件类型?以下是经过实践验证的事件协议设计:

事件类型描述典型字段前端处理方式
step_start步骤开始step_id, step_name, timestamp显示步骤指示器
step_data步骤中间数据step_id, content, data_type渲染中间结果
step_end步骤完成step_id, duration, result更新步骤状态为完成
tokenLLM tokencontent追加到文本缓冲
tool_call工具调用tool_name, args, call_id显示工具调用卡片
tool_result工具返回call_id, result, duration更新工具调用结果
thinkingAgent 思考过程content, depth显示折叠的思考链
error错误message, code, recoverable显示错误提示
interrupted被中断partial_output, reason保留部分输出
done完整完成full_output, metadata标记完成,保存结果

这种事件协议设计的关键在于结构化——每个事件都有明确的类型和字段,前端可以根据类型做出不同的渲染决策。例如 tool_call 和 tool_result 事件可以让前端渲染出工具调用的可视化卡片,而不是把工具信息混在 LLM 的文本输出中。这种结构化的事件系统也让前端可以实现更丰富的交互——比如点击工具调用卡片查看完整的输入参数和返回结果。

2.9 SSE 的性能优化技巧

在处理大量并发 SSE 连接时,有几个性能优化点值得关注:

1. 使用 uvicorn 的多 worker 模式

SSE 是长连接,单个 worker 能处理的并发连接数有限。使用 uvicorn --workers 4 启动多个 worker 进程可以提升并发能力。但注意每个 worker 是独立进程,不共享内存中的会话状态,需要使用 Redis 等外部存储来共享状态。

2. 使用异步 I/O,避免阻塞事件循环

SSE 连接会占用一个协程,如果在处理请求时使用了同步阻塞操作(如同步的数据库查询),会阻塞整个事件循环,影响所有 SSE 连接。所有 I/O 操作必须使用异步库(如 asyncpg 代替 psycopg2,aiohttp 代替 requests)。

3. 设置最大连接数限制

每个 SSE 连接都会占用内存和文件描述符。需要在应用层设置最大连接数限制,防止服务器资源耗尽:

from asyncio import Semaphore

# 全局连接数限制
MAX_CONNECTIONS = 1000
connection_semaphore = Semaphore(MAX_CONNECTIONS)

@app.post("/api/sse/chat")
async def sse_chat(request: Request):
    if connection_semaphore.locked():
        raise HTTPException(503, "服务器繁忙,请稍后重试")
    
    async with connection_semaphore:
        # 处理 SSE 流
        return StreamingResponse(...)

4. 使用 HTTP/2 减少 TCP 连接数

HTTP/2 支持多路复用,多个 SSE 流可以复用同一个 TCP 连接,减少了服务器端的文件描述符消耗。在 Nginx 中启用 HTTP/2 非常简单:listen 443 ssl http2;。

在这里插入图片描述


三、WebSocket 方案:双向通信与中间状态推送

3.1 WebSocket 与 SSE 的本质区别

SSE 是服务器→客户端的单向流,而 WebSocket 是全双工的双向通信。在 Agent 场景中,这意味着用户可以在 Agent 正在输出时,发送新的消息或控制指令——比如"停下来"、“换个方向”、“详细说说第三点”。

WebSocket 模式(双向)

双向实时

随时发消息

随时推数据

服务器

客户端

SSE 模式(单向)

数据流

新请求需新连接

服务器

客户端

3.2 何时选择 WebSocket

场景SSE 够用吗?推荐 WebSocket?
单轮对话流式输出✅ 完全够用不必要
多步骤任务执行✅ 基本够用可选
流式过程中用户可打断⚠️ 需额外连接✅ 推荐
多 Agent 实时协作❌ 不够✅ 推荐
Agent 主动推送通知❌ 不够✅ 推荐
实时多用户协作❌ 不够✅ 推荐
前端需发送大量中间数据❌ 不够✅ 推荐

3.3 后端实现:FastAPI + Python-SocketIO

import socketio
from openai import AsyncOpenAI

# 创建 Socket.IO 服务器
sio = socketio.AsyncServer(
    async_mode="asgi",
    cors_allowed_origins="*",
    ping_interval=20,   # 心跳间隔
    ping_timeout=30,   # 心跳超时
)

# LLM 客户端
llm_client = AsyncOpenAI(api_key="your-api-key", base_url="https://api.deepseek.com/v1")

# 存储每个会话的活跃生成任务
active_tasks: dict[str, asyncio.Task] = {}


@sio.event
async def connect(sid, environ):
    """客户端连接时触发"""
    print(f"客户端 {sid} 已连接")
    await sio.emit("connected", {"sid": sid}, to=sid)


@sio.event
async def disconnect(sid):
    """客户端断开时触发"""
    print(f"客户端 {sid} 已断开")
    # 如果有正在进行的生成任务,取消它
    if sid in active_tasks:
        active_tasks[sid].cancel()
        del active_tasks[sid]


@sio.on("chat")
async def handle_chat(sid, data):
    """
    处理用户消息,开始流式输出
    data: {"message": "用户消息", "session_id": "会话ID"}
    """
    user_message = data.get("message", "")
    
    # 推送"开始处理"事件
    await sio.emit("agent_event", {
        "type": "step",
        "step_id": "start",
        "status": "started",
        "message": "开始处理你的请求..."
    }, to=sid)
    
    # 创建异步任务处理 LLM 调用
    task = asyncio.create_task(stream_llm_response(sid, user_message))
    active_tasks[sid] = task


async def stream_llm_response(sid: str, user_message: str):
    """
    流式调用 LLM 并通过 WebSocket 推送结果
    """
    try:
        stream = await llm_client.chat.completions.create(
            model="deepseek-chat",
            messages=[
                {"role": "system", "content": "你是一个专业的AI助手。"},
                {"role": "user", "content": user_message},
            ],
            stream=True,
        )
        
        async for chunk in stream:
            # 检查任务是否被取消(用户中断)
            if sid not in active_tasks:
                print(f"会话 {sid} 已中断")
                return
            
            content = chunk.choices[0].delta.content
            if content:
                await sio.emit("agent_event", {
                    "type": "token",
                    "content": content
                }, to=sid)
        
        # 推送完成事件
        await sio.emit("agent_event", {
            "type": "done",
            "message": "回复完成"
        }, to=sid)
        
    except asyncio.CancelledError:
        # 用户主动中断
        await sio.emit("agent_event", {
            "type": "interrupted",
            "message": "生成已被中断"
        }, to=sid)
    except Exception as e:
        await sio.emit("agent_event", {
            "type": "error",
            "message": str(e)
        }, to=sid)
    finally:
        if sid in active_tasks:
            del active_tasks[sid]


@sio.on("interrupt")
async def handle_interrupt(sid, data):
    """
    用户主动中断当前生成
    """
    if sid in active_tasks:
        active_tasks[sid].cancel()
        # 注意:cancel() 会触发 CancelledError,由 stream_llm_response 内部处理

代码解析: 这段代码使用 Python-SocketIO 实现了 WebSocket 双向通信。与 SSE 的本质区别在于 active_tasks 字典——它维护了每个会话的活跃生成任务,使得用户可以通过 interrupt 事件随时中断正在进行的 LLM 生成。stream_llm_response 被封装为独立的异步任务,通过 asyncio.Task.cancel() 实现优雅中断。CancelledError 的捕获确保中断后能正确通知前端。ping_interval 和 ping_timeout 的配置防止了长连接被中间代理超时断开。

3.4 前端实现:Socket.IO 客户端

/**
 * WebSocket 流式 Agent 客户端
 * 支持双向通信:接收流式输出 + 发送中断指令
 */
import { io } from "socket.io-client";

class AgentWebSocketClient {
  constructor(url) {
    this.socket = io(url, {
      transports: ["websocket"],
      reconnection: true,
      reconnectionDelay: 1000,
      reconnectionAttempts: 5,
    });

    this.socket.on("connect", () => {
      console.log("已连接到 Agent 服务器");
    });

    this.socket.on("agent_event", (event) => {
      switch (event.type) {
        case "step":
          // 步骤状态更新
          this.onStep?.(event);
          break;
        case "token":
          // LLM token
          this.onToken?.(event.content);
          break;
        case "done":
          this.onDone?.();
          break;
        case "interrupted":
          this.onInterrupted?.();
          break;
        case "error":
          this.onError?.(event.message);
          break;
      }
    });

    this.socket.on("disconnect", () => {
      console.log("与 Agent 服务器断开连接");
    });
  }

  // 发送消息
  send(message) {
    this.socket.emit("chat", { message });
  }

  // 中断当前生成
  interrupt() {
    this.socket.emit("interrupt", {});
  }

  // 事件回调(由使用者设置)
  onToken = null;
  onStep = null;
  onDone = null;
  onInterrupted = null;
  onError = null;
}

// 使用示例
const agent = new AgentWebSocketClient("http://localhost:8000");
let outputText = "";

agent.onToken = (token) => {
  outputText += token;
  renderMarkdown(outputText);
};

agent.onStep = (event) => {
  updateStepIndicator(event.step_id, event.status, event.message);
};

agent.onDone = () => {
  setLoading(false);
};

agent.onInterrupted = () => {
  outputText += "\n\n[生成已被中断]";
  renderMarkdown(outputText);
  setLoading(false);
};

// 发送消息
agent.send("分析一下2024年AI Agent的发展趋势");

代码解析: 前端使用 Socket.IO 客户端库,它提供了比原生 WebSocket 更强大的能力——自动重连、消息确认、房间机制等。transports: ["websocket"] 强制使用 WebSocket 传输,避免降级到轮询。reconnection 配置确保网络抖动时自动恢复连接。回调函数的设计与 SSE 方案保持一致,使得切换底层协议时前端业务逻辑改动最小。最关键的新能力是 interrupt() 方法——用户可以在 Agent 输出过程中随时中断,这通过 WebSocket 的双向通信实现,无需新建连接。

3.5 WebSocket 的房间与命名空间机制

Socket.IO 提供了房间(Room)和命名空间(Namespace)机制,在多用户场景下非常有用。例如,多个用户可能同时与同一个 Agent 交互,或者一个 Agent 需要向多个用户广播状态:

@sio.on("join_room")
async def join_room(sid, data):
    """用户加入指定的 Agent 会话房间"""
    room = data.get("room_id")
    await sio.enter_room(sid, room)
    await sio.emit("room_joined", {"room_id": room}, to=sid)


@sio.on("agent_broadcast")
async def agent_broadcast(sid, data):
    """Agent 向房间内所有用户广播事件"""
    room = data.get("room_id")
    event = data.get("event")
    # 向房间内所有连接推送,包括发送者
    await sio.emit("agent_event", event, to=room)
    # 如果要排除发送者,使用 skip_sid=sid
    # await sio.emit("agent_event", event, room=room, skip_sid=sid)


# 使用命名空间隔离不同类型的 Agent
@sio.on("connect", namespace="/research-agent")
async def research_connect(sid, environ):
    print(f"研究 Agent 客户端 {sid} 已连接")
    await sio.emit("ready", {"agent_type": "research"}, to=sid, namespace="/research-agent")


@sio.on("chat", namespace="/research-agent")
async def research_chat(sid, data):
    """研究 Agent 专用聊天接口"""
    # 专门处理研究类查询
    message = data.get("message")
    # 启动研究 Agent 流式输出
    task = asyncio.create_task(stream_research_agent(sid, message))
    active_generations[f"{sid}:research"] = task

代码解析: 房间机制允许将多个连接分组——同一房间内的消息会广播给所有成员。这在多用户协作场景中非常有用,例如多个分析师同时查看同一 Agent 的研究报告生成过程。命名空间则用于隔离不同类型的 Agent 通信——研究 Agent 和代码 Agent 各有自己的命名空间,消息互不干扰。这种设计使得单个 Socket.IO 服务器可以同时服务多种 Agent,而不需要为每种 Agent 部署独立的服务器。

3.6 WebSocket 鉴权与安全

WebSocket 连接建立后的鉴权与 HTTP 不同——连接升级后无法再修改请求头。常见的 WebSocket 鉴权方案有三种:

鉴权方案实现优点缺点
URL 参数ws://host?token=xxx简单直接Token 暴露在日志中
Cookie依赖 HTTP Cookie自动携带,安全跨域问题
初始消息连接后先发 Auth 事件灵活增加一次 RTT

推荐使用 Cookie + 初始消息的双重验证方案:

@sio.event
async def connect(sid, environ):
    """连接时验证身份"""
    # 1. 检查 Cookie 中的会话 ID
    cookie_header = environ.get("HTTP_COOKIE", "")
    cookies = parse_cookie(cookie_header)
    session_id = cookies.get("session_id")
    
    if not session_id or not verify_session(session_id):
        # 拒绝连接
        return False  # 返回 False 会拒绝连接
    
    # 2. 将用户信息关联到会话
    user = get_user_from_session(session_id)
    await sio.save_session(sid, {"user_id": user.id, "name": user.name})
    
    await sio.emit("auth_success", {"user": user.name}, to=sid)


@sio.on("chat")
async def handle_chat(sid, data):
    # 从会话中获取用户信息
    session = await sio.get_session(sid)
    user_id = session["user_id"]
    print(f"用户 {session['name']} (sid={sid}) 发送消息")
    # ... 处理聊天

代码解析: WebSocket 鉴权的关键时机是 connect 事件——在连接建立之前进行身份验证,返回 False 可以拒绝连接。save_session 将用户信息绑定到 Socket.IO 会话,后续事件处理时可以通过 get_session 获取。这种方案结合了 Cookie 的自动携带特性和服务端会话验证,既安全又无需前端额外处理。

在这里插入图片描述


四、流式输出的前端渲染:打字机效果、Markdown 实时渲染

4.1 前端渲染的挑战

流式输出的前端渲染不是简单的"追加文本"这么容易。它面临几个关键挑战:

  1. Markdown 实时解析——LLM 输出的是 Markdown 格式文本,但流式接收时 Markdown 是不完整的。比如收到 ```python\nimport os\n 时,代码块还没结束,如何渲染?
  2. 性能——每个 token 都触发一次 DOM 更新会导致频繁重排。当输出达到数千字时,页面会明显卡顿。
  3. 滚动控制——用户可能在向上滚动查看历史内容,此时自动滚动到底部会打断阅读。
  4. 光标效果——打字机效果需要一个闪烁的光标指示当前输出位置。

4.2 打字机效果实现

/**
 * 打字机渲染器
 * 高性能流式文本渲染,支持光标效果和智能滚动
 */
class TypewriterRenderer {
  constructor(containerEl, options = {}) {
    this.container = containerEl;
    this.buffer = "";               // 完整文本缓冲
    this.renderedLength = 0;        // 已渲染长度
    this.cursor = options.cursor || "▋";
    this.cursorBlink = options.cursorBlink !== false;
    this.autoScroll = true;
    this.renderQueue = [];
    this.rendering = false;
    
    // 检测用户滚动
    this.container.addEventListener("scroll", () => {
      const { scrollTop, scrollHeight, clientHeight } = this.container;
      // 用户距离底部超过 50px 时,暂停自动滚动
      this.autoScroll = scrollHeight - scrollTop - clientHeight < 50;
    });
    
    // 光标闪烁
    if (this.cursorBlink) {
      setInterval(() => this.blinkCursor(), 500);
    }
  }

  /**
   * 追加 token 到缓冲区
   */
  append(token) {
    this.buffer += token;
    this.scheduleRender();
  }

  /**
   * 调度渲染(使用 requestAnimationFrame 合并多个 token)
   */
  scheduleRender() {
    if (this.rendering) return;
    this.rendering = true;
    requestAnimationFrame(() => {
      this.render();
      this.rendering = false;
    });
  }

  /**
   * 渲染当前缓冲区内容
   */
  render() {
    // 使用 markdown 解析库(如 marked.js)渲染
    const html = this.parseMarkdown(this.buffer);
    this.container.innerHTML = html + this.getCursorHTML();
    
    if (this.autoScroll) {
      this.container.scrollTop = this.container.scrollHeight;
    }
  }

  /**
   * Markdown 解析(处理不完整 Markdown)
   */
  parseMarkdown(text) {
    // 补全不完整的代码块
    const codeBlockCount = (text.match(/```/g) || []).length;
    if (codeBlockCount % 2 !== 0) {
      text = text + "\n```";  // 临时闭合代码块
    }
    
    // 使用 marked.js 解析
    if (window.marked) {
      return window.marked.parse(text);
    }
    // 降级为纯文本
    return text.replace(/\n/g, "<br>");
  }

  getCursorHTML() {
    return this.cursorVisible 
      ? `<span class="typing-cursor">${this.cursor}</span>` 
      : "";
  }

  blinkCursor() {
    this.cursorVisible = !this.cursorVisible;
    const cursor = this.container.querySelector(".typing-cursor");
    if (cursor) {
      cursor.style.opacity = this.cursorVisible ? "1" : "0";
    }
  }

  /**
   * 完成渲染(移除光标)
   */
  finish() {
    const html = this.parseMarkdown(this.buffer);
    this.container.innerHTML = html;
  }

  /**
   * 重置
   */
  reset() {
    this.buffer = "";
    this.renderedLength = 0;
    this.container.innerHTML = "";
  }
}

代码解析: 这个 TypewriterRenderer 解决了流式渲染的几个核心问题。第一,scheduleRender 使用 requestAnimationFrame 合并多个 token 的渲染请求——当 token 到达速度快于屏幕刷新率时,多余的渲染会被合并,避免性能浪费。第二,parseMarkdown 处理了不完整 Markdown 的问题——当检测到代码块标记 ````的数量为奇数时,临时补一个闭合标记,让 Markdown 解析器能正确渲染未完成的代码块。第三,scroll` 事件监听器实现了"智能滚动"——当用户向上滚动查看历史内容时,暂停自动滚到底部;当用户滚回底部附近时,恢复自动滚动。

4.3 React 组件实现

import React, { useState, useRef, useCallback, useEffect } from "react";
import { marked } from "marked";

/**
 * React 流式 Markdown 渲染组件
 * 支持打字机效果、代码高亮、智能滚动
 */
function StreamingMessage({ isStreaming }) {
  const [content, setContent] = useState("");
  const [isAtBottom, setIsAtBottom] = useState(true);
  const containerRef = useRef(null);
  const rafRef = useRef(null);

  // 智能滚动检测
  const handleScroll = useCallback(() => {
    const el = containerRef.current;
    if (!el) return;
    const { scrollTop, scrollHeight, clientHeight } = el;
    setIsAtBottom(scrollHeight - scrollTop - clientHeight < 50);
  }, []);

  // 追加 token(暴露给父组件使用)
  const appendToken = useCallback((token) => {
    setContent((prev) => {
      const newContent = prev + token;
      // 使用 rAF 合并 DOM 更新
      if (rafRef.current) cancelAnimationFrame(rafRef.current);
      rafRef.current = requestAnimationFrame(() => {
        if (isAtBottom && containerRef.current) {
          containerRef.current.scrollTop = containerRef.current.scrollHeight;
        }
      });
      return newContent;
    });
  }, [isAtBottom]);

  // 暴露 appendToken 给父组件
  useImperativeHandle(ref, () => ({ appendToken, reset: () => setContent("") }));

  // 渲染 Markdown
  const html = useMemo(() => {
    let text = content;
    // 补全不完整的代码块
    const fenceCount = (text.match(/```/g) || []).length;
    if (fenceCount % 2 !== 0) text += "\n```";
    return marked.parse(text);
  }, [content]);

  return (
    <div
      ref={containerRef}
      onScroll={handleScroll}
      className="streaming-message"
      dangerouslySetInnerHTML={{
        __html: html + (isStreaming ? '<span class="cursor">▋</span>' : "")
      }}
    />
  );
}

代码解析: 这是 React 版本的流式渲染组件。useMemo 对 Markdown 解析做了缓存优化——只有 content 变化时才重新解析。useImperativeHandle 将 appendToken 方法暴露给父组件,使得 SSE/WebSocket 客户端可以直接调用它追加 token。dangerouslySetInnerHTML 看起来危险,但因为内容经过 marked.parse 转义处理,XSS 风险已被消除。光标元素在流式结束后自动移除。配合 CSS 动画 @keyframes blink,光标会闪烁。

4.3 代码高亮的流式处理

LLM 输出经常包含代码块,而流式接收时代码块是不完整的。如何在代码还在生成时就开始语法高亮?这里的关键是使用支持增量更新的代码高亮库,如 highlight.js 或 Prism.js 的实时模式:

/**
 * 流式代码高亮处理器
 * 在代码块生成过程中持续更新高亮
 */
class StreamingCodeHighlighter {
  constructor() {
    this.codeBlockCount = 0;
    this.currentBlock = null;
  }

  /**
   * 处理流式文本
   * 检测代码块边界,在代码块结束时触发高亮
   */
  process(text) {
    // 统计代码块标记数量
    const fenceCount = (text.match(/```/g) || []).length;
    
    if (fenceCount % 2 === 0) {
      // 所有代码块都完整,可以安全高亮
      return this.highlightAll(text);
    } else {
      // 有未闭合的代码块
      // 临时闭合后高亮,但标记最后一个代码块为"正在输入"
      const tempText = text + "\n```";
      const highlighted = this.highlightAll(tempText);
      
      // 给最后一个代码块添加"正在输入"样式
      return highlighted.replace(
        /(<pre[^>]*>)([\s\S]*?)(<\/pre>)(?!.*<pre)/,
        '$1$2<span class="typing-cursor">▋</span>$3'
      );
    }
  }

  highlightAll(text) {
    // 使用 marked.js 解析 Markdown
    const html = marked.parse(text);
    
    // 使用 highlight.js 对 <pre><code> 块进行高亮
    const wrapper = document.createElement("div");
    wrapper.innerHTML = html;
    wrapper.querySelectorAll("pre code").forEach((block) => {
      hljs.highlightElement(block);
    });
    
    return wrapper.innerHTML;
  }
}

代码解析: 代码高亮的核心挑战是代码块未闭合时的语法解析问题。highlight.js 需要完整的代码才能正确高亮,但流式输出时代码是不完整的。解决方案是检测 ` ````标记的数量——偶数表示所有代码块完整,奇数表示有未闭合的代码块。对于未闭合的情况,临时添加一个闭合标记让 Markdown 解析器能正常工作,然后在渲染后的代码块末尾添加闪烁光标。每次新 token 到达时都会重新高亮整个代码块——虽然有些性能开销,但由于代码块通常不长,现代浏览器可以轻松应对。如果性能确实是问题,可以加入节流逻辑,每 100ms 更新一次代码高亮。

4.4 渲染性能优化实践

当 Agent 的输出达到数千字时(如生成长报告),前端渲染会面临性能瓶颈。以下是经过实践验证的优化策略:

策略一:分层渲染

将已完成的段落和正在生成的段落分离——已完成的段落不再重新渲染,只更新正在生成的最后一个段落:

class LayeredRenderer {
  constructor(container) {
    this.container = container;
    this.completedParagraphs = [];  // 已完成的段落
    this.currentParagraph = "";      // 正在生成的段落
  }

  append(token) {
    this.currentParagraph += token;
    
    // 检测段落结束(双换行)
    if (this.currentParagraph.includes("\n\n")) {
      const parts = this.currentParagraph.split("\n\n");
      // 已完成的段落移到 completedParagraphs
      this.completedParagraphs.push(parts.shift());
      this.currentParagraph = parts.join("\n\n");
    }
    
    this.render();
  }

  render() {
    // 只渲染容器,已完成的段落用 innerHTML 设置一次
    const completedHtml = this.completedParagraphs
      .map(p => marked.parse(p))
      .join("");
    const currentHtml = marked.parse(this.currentParagraph);
    
    this.container.innerHTML = 
      completedHtml + 
      '<div class="streaming-paragraph">' + currentHtml + 
      '<span class="cursor">▋</span></div>';
  }
}

策略二:虚拟滚动

对于超长输出(如万字报告),可以使用虚拟滚动技术——只渲染视口可见的 DOM 元素。这需要使用 react-window 或 vue-virtual-scroller 等库。虽然实现复杂度增加,但可以支持无限长度的流式输出而不卡顿。

策略三:Web Worker 解析

将 Markdown 解析和代码高亮放到 Web Worker 中执行,避免阻塞主线程。主线程只负责接收 token 和更新 DOM。这在处理包含大量代码块的输出时效果尤为明显。


五、中断处理:用户中途打断 Agent 的处理策略

5.1 中断的场景与复杂度

Agent 的中断处理远比"停止生成"复杂。一个完整的 Agent 任务可能涉及多个资源:LLM 流式生成、外部工具调用、数据库写入、文件创建等。中断时需要确保:

  1. LLM 生成停止——取消正在进行的流式请求
  2. 部分结果保留——已生成的内容不丢失,用户可以参考
  3. 资源清理——已打开的文件、数据库连接等正确关闭
  4. 状态一致性——Agent 的内部状态保持一致,不留下"半成品"状态
  5. 通知前端——前端收到中断确认,更新 UI

用户发送消息

LLM 生成中

用户点击停止

工具调用中

用户点击停止

工具返回,继续生成

生成完成

清理资源

返回结果

超时清理

Idle

Running

Generating

Interrupted

ToolCalling

Completed

中断时需要:
1. 取消 LLM 请求
2. 保留部分输出
3. 清理资源
4. 通知前端

5.2 后端中断处理实现

import asyncio
from dataclasses import dataclass, field
from enum import Enum
from typing import Optional, Any


class AgentState(Enum):
    IDLE = "idle"
    RUNNING = "running"
    GENERATING = "generating"
    TOOL_CALLING = "tool_calling"
    INTERRUPTED = "interrupted"
    COMPLETED = "completed"
    FAILED = "failed"


@dataclass
class AgentSession:
    """Agent 会话上下文,管理状态和资源"""
    session_id: str
    state: AgentState = AgentState.IDLE
    llm_task: Optional[asyncio.Task] = None
    tool_tasks: list[asyncio.Task] = field(default_factory=list)
    partial_output: str = ""
    cleanup_callbacks: list = field(default_factory=list)
    
    async def interrupt(self):
        """
        中断当前 Agent 执行
        按顺序取消所有任务并清理资源
        """
        if self.state in (AgentState.COMPLETED, AgentState.INTERRUPTED):
            return
        
        self.state = AgentState.INTERRUPTED
        print(f"[{self.session_id}] 开始中断 Agent...")
        
        # 1. 取消 LLM 生成任务
        if self.llm_task and not self.llm_task.done():
            self.llm_task.cancel()
            try:
                await asyncio.wait_for(self.llm_task, timeout=2.0)
            except (asyncio.CancelledError, asyncio.TimeoutError):
                pass
        
        # 2. 取消工具调用任务
        for task in self.tool_tasks:
            if not task.done():
                task.cancel()
                try:
                    await asyncio.wait_for(task, timeout=2.0)
                except (asyncio.CancelledError, asyncio.TimeoutError):
                    pass
        
        # 3. 执行清理回调(关闭文件、数据库连接等)
        for callback in self.cleanup_callbacks:
            try:
                result = callback()
                if asyncio.iscoroutine(result):
                    await result
            except Exception as e:
                print(f"清理回调出错: {e}")
        
        # 4. 保留部分输出
        print(f"[{self.session_id}] 中断完成,已保留 {len(self.partial_output)} 字符的部分输出")
        
        self.state = AgentState.IDLE


class InterruptibleAgent:
    """
    支持中断的 Agent
    """
    
    def __init__(self):
        self.sessions: dict[str, AgentSession] = {}
    
    async def run(self, session_id: str, user_message: str, emit_callback):
        """
        运行 Agent,支持被中断
        emit_callback: 用于推送事件的回调函数
        """
        session = AgentSession(session_id=session_id)
        self.sessions[session_id] = session
        
        try:
            session.state = AgentState.RUNNING
            await emit_callback({"type": "state", "state": "running"})
            
            # 注册清理函数
            session.cleanup_callbacks.append(
                lambda: print(f"[{session_id}] 清理: 关闭数据库连接")
            )
            session.cleanup_callbacks.append(
                lambda: print(f"[{session_id}] 清理: 删除临时文件")
            )
            
            # LLM 流式生成
            session.state = AgentState.GENERATING
            await emit_callback({"type": "state", "state": "generating"})
            
            async for token in self._stream_llm(user_message, session):
                if session.state == AgentState.INTERRUPTED:
                    # 被中断,停止生成
                    await emit_callback({
                        "type": "interrupted",
                        "partial_output": session.partial_output,
                        "message": "Agent 已被中断,以上为部分输出"
                    })
                    return
                
                session.partial_output += token
                await emit_callback({"type": "token", "content": token})
            
            session.state = AgentState.COMPLETED
            await emit_callback({"type": "done", "output": session.partial_output})
            
        except asyncio.CancelledError:
            await session.interrupt()
            raise
        except Exception as e:
            session.state = AgentState.FAILED
            await emit_callback({"type": "error", "message": str(e)})
        finally:
            if session_id in self.sessions:
                del self.sessions[session_id]
    
    async def _stream_llm(self, message: str, session: AgentSession):
        """模拟 LLM 流式输出(实际应调用 LLM API)"""
        response = "这是一个模拟的流式响应,用于演示中断机制。每个字都是单独生成的。"
        for char in response:
            if session.state == AgentState.INTERRUPTED:
                return
            yield char
            await asyncio.sleep(0.05)  # 模拟生成延迟
    
    async def interrupt(self, session_id: str):
        """外部调用:中断指定会话"""
        if session_id in self.sessions:
            await self.sessions[session_id].interrupt()

代码解析: 这段代码实现了一个完整的可中断 Agent 框架。AgentSession 是核心数据结构,它跟踪 Agent 的状态(通过状态机管理)、持有所有活跃任务的引用、保存部分输出、维护清理回调列表。interrupt() 方法的清理顺序经过精心设计:先取消 LLM 任务(最耗资源),再取消工具任务,然后执行资源清理回调,最后保留部分输出。每个取消操作都有 2 秒超时保护,防止任务无法被取消时无限等待。cleanup_callbacks 的设计允许各步骤注册自己的清理逻辑(如关闭文件、回滚事务),实现了关注点分离。

5.3 前端中断 UI 设计

/**
 * 带中断功能的 Agent 对话组件
 */
function AgentChat() {
  const [messages, setMessages] = useState([]);
  const [isStreaming, setIsStreaming] = useState(false);
  const [currentMessage, setCurrentMessage] = useState("");
  const clientRef = useRef(null);

  const handleSend = async (text) => {
    setIsStreaming(true);
    setCurrentMessage("");

    // 根据 URL 选择 SSE 或 WebSocket 客户端
    const client = new AgentWebSocketClient("ws://localhost:8000");
    clientRef.current = client;

    client.onToken = (token) => {
      setCurrentMessage((prev) => prev + token);
    };

    client.onDone = () => {
      setMessages((prev) => [...prev, { role: "agent", content: currentMessageRef.current }]);
      setCurrentMessage("");
      setIsStreaming(false);
    };

    client.onInterrupted = () => {
      // 保留部分输出
      setMessages((prev) => [
        ...prev,
        { role: "agent", content: currentMessageRef.current + "\n\n[已中断]" },
      ]);
      setCurrentMessage("");
      setIsStreaming(false);
    };

    client.send(text);
  };

  const handleStop = () => {
    if (clientRef.current) {
      clientRef.current.interrupt();
    }
  };

  return (
    <div className="chat-container">
      {/* 消息列表 */}
      <div className="messages">
        {messages.map((msg, i) => (
          <div key={i} className={`message ${msg.role}`}>
            <ReactMarkdown>{msg.content}</ReactMarkdown>
          </div>
        ))}
        {/* 正在生成的消息 */}
        {isStreaming && currentMessage && (
          <div className="message agent streaming">
            <ReactMarkdown>{currentMessage}</ReactMarkdown>
            <span className="cursor">▋</span>
          </div>
        )}
      </div>

      {/* 输入区域 */}
      <div className="input-area">
        {isStreaming ? (
          <button onClick={handleStop} className="stop-btn">
            ⏹ 停止生成
          </button>
        ) : (
          <ChatInput onSend={handleSend} />
        )}
      </div>
    </div>
  );
}

代码解析: 前端中断 UI 的关键是状态管理。isStreaming 控制输入区域和停止按钮的切换——当 Agent 正在生成时,输入框变成停止按钮,防止用户发送新消息。onInterrupted 回调保留了部分输出并标记 [已中断],让用户知道这不是完整回答。这种设计让用户始终掌控 Agent 的行为,增强了信任感。

5.3 中断后的恢复策略

中断并不意味着放弃。好的 Agent 设计应该支持中断后的恢复——用户可以选择从断点继续生成,而不是从头开始。这需要后端保存中断时的上下文:

@dataclass
class InterruptedSession:
    """中断会话的上下文保存"""
    session_id: str
    original_message: str          # 原始用户消息
    conversation_history: list      # 对话历史
    partial_output: str             # 已生成的部分输出
    interrupt_reason: str           # 中断原因(用户主动、超时等)
    timestamp: float                # 中断时间
    llm_context: dict = None        # LLM 上下文(如缓存的 KV Cache)

    def resume_prompt(self) -> str:
    	"""生成恢复提示词"""
        return (
            f"之前的回复被中断,已生成的部分内容如下:\n\n"
            f"{self.partial_output}\n\n"
            f"请从断点处继续,不要重复已有内容。"
        )


# 存储中断的会话(实际应使用 Redis)
interrupted_sessions: dict[str, InterruptedSession] = {}


async def handle_resume(sid: str, session_id: str):
    """从中断点恢复生成"""
    session = interrupted_sessions.get(session_id)
    if not session:
        await sio.emit("event", {
            "type": "error",
            "message": "中断会话已过期,无法恢复"
        }, to=sid)
        return
    
    # 使用部分输出作为上下文,继续生成
    messages = session.conversation_history + [
        {"role": "assistant", "content": session.partial_output},
        {"role": "user", "content": session.resume_prompt()}
    ]
    
    await sio.emit("event", {
        "type": "step",
        "step": "resuming",
        "status": "started",
        "message": f"从断点恢复,已有 {len(session.partial_output)} 字符"
    }, to=sid)
    
    # 继续流式生成
    task = asyncio.create_task(_stream_to_client(sid, session.resume_prompt(), messages))
    active_generations[sid] = task

代码解析: 中断恢复的关键是保存足够多的上下文信息。InterruptedSession 不仅保存了已生成的部分输出,还保存了完整的对话历史和原始用户消息,使得恢复时可以重建完整的 LLM 上下文。resume_prompt 方法构造了一个特殊的提示,告诉 LLM 从断点继续而不是重新开始。需要注意的是,由于 LLM 的无状态特性,恢复后的生成质量可能不如一次性生成——因为 LLM 需要根据部分输出推断之前的思路。在生产环境中,可以考虑保存 LLM 的 KV Cache(如果模型 API 支持),这样可以实现真正的断点续传。

5.4 不同中断场景的处理策略

Agent 的中断可能由多种原因触发,不同原因需要不同的处理策略:

中断原因触发方式处理策略恢复选项
用户主动停止点击停止按钮立即取消任务,保留部分输出支持恢复生成
网络断开连接超时/断开保存上下文到 Redis,等待重连自动恢复
服务器错误异常/崩溃记录错误日志,清理资源手动重试
超时生成时间过长软中断,通知用户等待过长可继续等待或停止
并发限制超过最大连接数排队等待可自动开始

六、实战:一个完整的流式 Agent 前后端实现

6.1 项目结构

streaming-agent/
├── backend/
│   ├── main.py              # FastAPI 入口,挂载 SSE + WebSocket
│   ├── sse_handler.py        # SSE 流式处理
│   ├── ws_handler.py        # WebSocket 流式处理
│   ├── agent.py             # Agent 核心逻辑
│   └── interrupt.py         # 中断管理
├── frontend/
│   ├── index.html           # 页面入口
│   ├── app.js              # 前端逻辑
│   └── styles.css          # 样式
├── requirements.txt
└── README.md

6.2 完整后端入口

"""
流式 Agent 完整后端
同时提供 SSE 和 WebSocket 两种接口
"""
import os
import json
import asyncio
from fastapi import FastAPI, Request, HTTPException
from fastapi.responses import StreamingResponse
from fastapi.middleware.cors import CORSMiddleware
from openai import AsyncOpenAI

app = FastAPI(title="Streaming Agent")
app.add_middleware(
    CORSMiddleware,
    allow_origins=["*"],
    allow_methods=["*"],
    allow_headers=["*"],
)

llm_client = AsyncOpenAI(
    api_key=os.getenv("LLM_API_KEY", "sk-xxx"),
    base_url=os.getenv("LLM_BASE_URL", "https://api.deepseek.com/v1")
)

# ====== SSE 接口 ======

@app.post("/api/sse/chat")
async def sse_chat(request: Request):
    """SSE 流式对话接口"""
    body = await request.json()
    message = body.get("message", "")
    conversation = body.get("conversation", [])
    
    async def event_stream():
        try:
            # 推送步骤:思考中
            yield format_sse({"type": "step", "step": "thinking", "status": "started", "message": "正在理解你的问题..."})
            await asyncio.sleep(0.3)
            yield format_sse({"type": "step", "step": "thinking", "status": "completed"})
            
            # 推送步骤:生成中
            yield format_sse({"type": "step", "step": "generating", "status": "started", "message": "正在生成回复..."})
            
            # 调用 LLM
            messages = [
                {"role": "system", "content": "你是专业AI助手,回答要准确、简洁。"}
            ] + conversation + [{"role": "user", "content": message}]
            
            stream = await llm_client.chat.completions.create(
                model="deepseek-chat",
                messages=messages,
                stream=True,
                max_tokens=2000,
            )
            
            full_text = ""
            async for chunk in stream:
                if await request.is_disconnected():
                    yield format_sse({"type": "interrupted", "partial": full_text})
                    return
                
                token = chunk.choices[0].delta.content
                if token:
                    full_text += token
                    yield format_sse({"type": "token", "content": token})
            
            yield format_sse({"type": "step", "step": "generating", "status": "completed"})
            yield format_sse({"type": "done", "full_text": full_text})
            
        except Exception as e:
            yield format_sse({"type": "error", "message": str(e)})
    
    return StreamingResponse(
        event_stream(),
        media_type="text/event-stream",
        headers={"Cache-Control": "no-cache", "X-Accel-Buffering": "no"}
    )


def format_sse(data: dict) -> str:
    return f"data: {json.dumps(data, ensure_ascii=False)}\n\n"


# ====== WebSocket 接口 ======

import socketio

sio = socketio.AsyncServer(async_mode="asgi", cors_allowed_origins="*")
socket_app = socketio.ASGIApp(sio)

active_generations: dict[str, asyncio.Task] = {}


@sio.on("chat")
async def ws_chat(sid, data):
    message = data.get("message", "")
    conversation = data.get("conversation", [])
    
    await sio.emit("event", {"type": "step", "step": "thinking", "status": "started"}, to=sid)
    await asyncio.sleep(0.3)
    await sio.emit("event", {"type": "step", "step": "thinking", "status": "completed"}, to=sid)
    await sio.emit("event", {"type": "step", "step": "generating", "status": "started"}, to=sid)
    
    task = asyncio.create_task(_stream_to_client(sid, message, conversation))
    active_generations[sid] = task


async def _stream_to_client(sid: str, message: str, conversation: list):
    try:
        messages = [
            {"role": "system", "content": "你是专业AI助手,回答要准确、简洁。"}
        ] + conversation + [{"role": "user", "content": message}]
        
        stream = await llm_client.chat.completions.create(
            model="deepseek-chat",
            messages=messages,
            stream=True,
            max_tokens=2000,
        )
        
        async for chunk in stream:
            if sid not in active_generations:
                await sio.emit("event", {"type": "interrupted"}, to=sid)
                return
            
            token = chunk.choices[0].delta.content
            if token:
                await sio.emit("event", {"type": "token", "content": token}, to=sid)
        
        await sio.emit("event", {"type": "done"}, to=sid)
        
    except asyncio.CancelledError:
        await sio.emit("event", {"type": "interrupted"}, to=sid)
    except Exception as e:
        await sio.emit("event", {"type": "error", "message": str(e)}, to=sid)
    finally:
        active_generations.pop(sid, None)


@sio.on("interrupt")
async def ws_interrupt(sid, data):
    if sid in active_generations:
        active_generations[sid].cancel()


# 挂载 Socket.IO
app.mount("/ws", socket_app)


if __name__ == "__main__":
    import uvicorn
    uvicorn.run(app, host="0.0.0.0", port=8000)

代码解析: 这是一个同时提供 SSE 和 WebSocket 两种接口的完整后端。SSE 接口挂载在 /api/sse/chat,适合简单的对话场景;WebSocket 接口挂载在 /ws 路径下,适合需要中断控制的场景。两个接口共享同一个 LLM 客户端和对话逻辑,只是传输机制不同。format_sse 工具函数确保所有 SSE 数据格式一致。active_generations 字典管理 WebSocket 会话的生成任务,支持用户随时中断。app.mount("/ws", socket_app) 将 Socket.IO 应用挂载为子应用,与 FastAPI 的路由共存。

6.3 完整前端实现

<!DOCTYPE html>
<html lang="zh-CN">
<head>
    <meta charset="UTF-8">
    <meta name="viewport" content="width=device-width, initial-scale=1.0">
    <title>流式 Agent Demo</title>
    <style>
        * { margin: 0; padding: 0; box-sizing: border-box; }
        body { font-family: -apple-system, sans-serif; background: #1a1a2e; color: #eee; }
        .container { max-width: 800px; margin: 0 auto; height: 100vh; display: flex; flex-direction: column; }
        .header { padding: 16px; border-bottom: 1px solid #333; display: flex; justify-content: space-between; align-items: center; }
        .messages { flex: 1; overflow-y: auto; padding: 16px; }
        .msg { margin-bottom: 16px; padding: 12px 16px; border-radius: 12px; max-width: 80%; }
        .msg.user { background: #16213e; margin-left: auto; }
        .msg.agent { background: #0f3460; }
        .msg.agent.streaming .cursor { animation: blink 0.8s infinite; }
        @keyframes blink { 0%, 50% { opacity: 1; } 51%, 100% { opacity: 0; } }
        .input-area { padding: 16px; border-top: 1px solid #333; display: flex; gap: 8px; }
        .input-area input { flex: 1; padding: 12px; border: 1px solid #333; border-radius: 8px; background: #16213e; color: #eee; }
        .btn { padding: 12px 24px; border: none; border-radius: 8px; cursor: pointer; font-size: 14px; }
        .btn-send { background: #e94560; color: white; }
        .btn-stop { background: #e74c3c; color: white; }
        .status-bar { padding: 4px 16px; font-size: 12px; color: #888; min-height: 20px; }
        .step-indicator { display: inline-block; margin-right: 8px; padding: 2px 8px; border-radius: 4px; background: #16213e; font-size: 11px; }
        .step-indicator.active { background: #e94560; color: white; }
        .step-indicator.done { background: #2ecc71; color: white; }
    </style>
</head>
<body>
    <div class="container">
        <div class="header">
            <h1>🤖 流式 Agent</h1>
            <select id="protocol">
                <option value="sse">SSE 模式</option>
                <option value="ws">WebSocket 模式</option>
            </select>
        </div>
        
        <div class="status-bar" id="status-bar"></div>
        
        <div class="messages" id="messages"></div>
        
        <div class="input-area">
            <input type="text" id="input" placeholder="输入消息..." />
            <button class="btn btn-send" id="send-btn" onclick="send()">发送</button>
            <button class="btn btn-stop" id="stop-btn" onclick="stop()" style="display:none">停止</button>
        </div>
    </div>

    <script>
        let isStreaming = false;
        let currentController = null;
        let wsSocket = null;
        const messagesEl = document.getElementById("messages");
        const inputEl = document.getElementById("input");
        const statusEl = document.getElementById("status-bar");
        const protocolEl = document.getElementById("protocol");
        const sendBtn = document.getElementById("send-btn");
        const stopBtn = document.getElementById("stop-btn");

        inputEl.addEventListener("keypress", (e) => {
            if (e.key === "Enter" && !isStreaming) send();
        });

        async function send() {
            const message = inputEl.value.trim();
            if (!message || isStreaming) return;
            
            inputEl.value = "";
            isStreaming = true;
            sendBtn.style.display = "none";
            stopBtn.style.display = "block";
            
            // 渲染用户消息
            appendMessage("user", message);
            
            // 创建 Agent 消息容器
            const agentMsg = appendMessage("agent", "");
            agentMsg.classList.add("streaming");
            
            const protocol = protocolEl.value;
            if (protocol === "sse") {
                await sendSSE(message, agentMsg);
            } else {
                sendWebSocket(message, agentMsg);
            }
        }

        async function sendSSE(message, agentMsg) {
            currentController = new AbortController();
            let fullText = "";
            
            try {
                const response = await fetch("/api/sse/chat", {
                    method: "POST",
                    headers: { "Content-Type": "application/json" },
                    body: JSON.stringify({ message }),
                    signal: currentController.signal,
                });
                
                const reader = response.body.getReader();
                const decoder = new TextDecoder();
                let buffer = "";
                
                while (true) {
                    const { done, value } = await reader.read();
                    if (done) break;
                    
                    buffer += decoder.decode(value, { stream: true });
                    const lines = buffer.split("\n\n");
                    buffer = lines.pop();
                    
                    for (const line of lines) {
                        if (!line.startsWith("data: ")) continue;
                        const data = line.slice(6).trim();
                        if (data === "[DONE]") continue;
                        
                        const event = JSON.parse(data);
                        handleEvent(event, agentMsg, (t) => {
                            fullText += t;
                            agentMsg.innerHTML = renderMarkdown(fullText) + '<span class="cursor">▋</span>';
                            scrollToBottom();
                        });
                    }
                }
                
                agentMsg.classList.remove("streaming");
                agentMsg.innerHTML = renderMarkdown(fullText);
                
            } catch (err) {
                if (err.name === "AbortError") {
                    agentMsg.innerHTML = renderMarkdown(fullText) + '<p style="color:#888">[已中断]</p>';
                } else {
                    agentMsg.innerHTML = `<p style="color:#e74c3c">错误: ${err.message}</p>`;
                }
            } finally {
                resetState();
            }
        }

        function sendWebSocket(message, agentMsg) {
            let fullText = "";
            
            wsSocket = io("ws://localhost:8000/ws", { transports: ["websocket"] });
            
            wsSocket.on("event", (event) => {
                handleEvent(event, agentMsg, (t) => {
                    fullText += t;
                    agentMsg.innerHTML = renderMarkdown(fullText) + '<span class="cursor">▋</span>';
                    scrollToBottom();
                });
                
                if (event.type === "done" || event.type === "interrupted" || event.type === "error") {
                    agentMsg.classList.remove("streaming");
                    if (event.type === "interrupted") {
                        agentMsg.innerHTML = renderMarkdown(fullText) + '<p style="color:#888">[已中断]</p>';
                    }
                    wsSocket.disconnect();
                    resetState();
                }
            });
            
            wsSocket.emit("chat", { message });
        }

        function handleEvent(event, agentMsg, onToken) {
            switch (event.type) {
                case "step":
                    updateStatus(event);
                    break;
                case "token":
                    onToken(event.content);
                    break;
                case "done":
                    updateStatus({ type: "step", step: "completed", status: "done", message: "完成" });
                    break;
                case "interrupted":
                    updateStatus({ message: "已中断" });
                    break;
                case "error":
                    updateStatus({ message: "错误: " + event.message });
                    break;
            }
        }

        function updateStatus(event) {
            if (event.message) {
                statusEl.textContent = event.message;
            }
        }

        function stop() {
            if (protocolEl.value === "sse") {
                currentController?.abort();
            } else if (wsSocket) {
                wsSocket.emit("interrupt", {});
            }
        }

        function appendMessage(role, content) {
            const div = document.createElement("div");
            div.className = `msg ${role}`;
            div.innerHTML = role === "user" ? escapeHtml(content) : renderMarkdown(content);
            messagesEl.appendChild(div);
            scrollToBottom();
            return div;
        }

        function renderMarkdown(text) {
            // 简易 Markdown 渲染(生产环境请使用 marked.js)
            return text
                .replace(/&/g, "&amp;")
                .replace(/</g, "&lt;")
                .replace(/>/g, "&gt;")
                .replace(/```(\w*)\n([\s\S]*?)```/g, '<pre><code>$2</code></pre>')
                .replace(/`([^`]+)`/g, '<code>$1</code>')
                .replace(/\*\*(.+?)\*\*/g, '<strong>$1</strong>')
                .replace(/\n/g, '<br>');
        }

        function escapeHtml(text) {
            const div = document.createElement("div");
            div.textContent = text;
            return div.innerHTML;
        }

        function scrollToBottom() {
            messagesEl.scrollTop = messagesEl.scrollHeight;
        }

        function resetState() {
            isStreaming = false;
            sendBtn.style.display = "block";
            stopBtn.style.display = "none";
            statusEl.textContent = "";
        }
    </script>
</body>
</html>

代码解析: 这是一个完整的单页面前端实现,同时支持 SSE 和 WebSocket 两种模式,用户可以通过下拉框切换。页面包含消息列表、状态栏、输入区域和停止按钮。SSE 模式使用 fetch + ReadableStream 接收流式数据,通过 AbortController 实现中断;WebSocket 模式使用 Socket.IO 客户端,通过 emit("interrupt") 事件实现中断。renderMarkdown 函数提供了简易的 Markdown 渲染能力(生产环境建议使用 marked.js),支持代码块、行内代码、粗体和换行。光标闪烁效果通过 CSS @keyframes blink 动画实现。状态栏实时显示 Agent 当前步骤(思考中→生成中→完成),让用户始终了解 Agent 的工作状态。整个页面无任何框架依赖,原生 HTML + JS + CSS 即可运行。

6.4 完整的前端渲染效果

当用户发送消息后,界面会经历以下状态变化:

  1. 发送瞬间:用户消息显示在右侧(蓝色背景),下方出现一个空的 Agent 消息气泡(深蓝色背景),带闪烁光标
  2. 思考阶段:状态栏显示"正在理解你的问题…",Agent 气泡为空
  3. 生成阶段:状态栏显示"正在生成回复…",文字逐字出现在 Agent 气泡中
  4. 完成阶段:光标消失,状态栏显示"完成",输入框恢复可用
  5. 中断场景:用户点击"停止"按钮,生成立即停止,已生成的文字保留,下方标注"[已中断]"

在这里插入图片描述


七、适用边界与风险提示

7.1 SSE 的适用边界

SSE 方案虽然简洁,但有其局限性:

  • 单向通信:只支持服务器→客户端,用户中断需要额外的 HTTP 请求或 WebSocket 连接
  • 连接数限制:浏览器对同一域名的 SSE 连接数有限制(HTTP/1.1 下通常 6 个),HTTP/2 下有所改善
  • POST 请求不支持:原生 EventSource API 只支持 GET,需要用 fetch 替代
  • 代理兼容性:某些反向代理(如旧版 Nginx)可能缓冲 SSE 流,需要配置 X-Accel-Buffering: no
  • 移动端不稳定:iOS Safari 对 SSE 的后台连接有较严格的超时限制

SSE 适合的场景:

  • 单轮对话流式输出
  • 简单的步骤状态推送
  • 不需要用户在流式过程中发送消息的场景
  • 快速原型开发

7.2 WebSocket 的适用边界

WebSocket 功能强大,但复杂度也更高:

  • 协议升级开销:连接建立时需要 HTTP→WS 协议升级,略增加首次连接延迟
  • 运维复杂度高:长连接管理、心跳机制、连接状态同步都需要额外处理
  • 负载均衡挑战:WebSocket 是有状态连接,负载均衡需要使用 IP Hash 或 Sticky Session
  • 移动端兼容性:部分移动网络(如某些 4G 代理)可能不支持 WebSocket 升级
  • 安全考量:需要额外实现鉴权、消息校验,防止未授权访问和注入攻击

WebSocket 适合的场景:

  • 需要双向实时通信的复杂 Agent
  • 多用户协作场景
  • Agent 需要主动推送通知
  • 对中断实时性要求高的场景

7.3 风险提示与生产环境建议

风险点影响建议
连接泄漏服务器资源耗尽设置连接超时 + 心跳检测 + 最大连接数限制
消息洪泛前端渲染卡顿使用 requestAnimationFrame 合并渲染 + 限流
XSS 注入安全漏洞对 LLM 输出做 HTML 转义 + CSP 策略
中断后资源未清理内存/连接泄漏使用 try/finally 确保资源释放 + 超时兜底
网络抖动导致断连用户体验中断实现自动重连 + 断点续传
并发连接过多服务器压力大实现连接池 + 排队机制
LLM 输出不当内容合规风险添加输出过滤层 + 敏感词检测

7.4 生产环境检查清单

上线前,请逐项检查以下内容:

  • SSE/Ws 连接设置了合理的超时时间(建议 5-10 分钟)
  • 实现了心跳检测,能及时发现断连
  • 前端实现了自动重连机制(指数退避重试)
  • 服务器设置了最大连接数限制
  • LLM 输出经过 HTML 转义,防止 XSS
  • 中断后所有资源(Task、连接、文件)被正确清理
  • 负载均衡配置支持长连接(IP Hash 或 Sticky Session)
  • Nginx 配置关闭了 SSE 响应缓冲
  • 监控告警覆盖了连接数、延迟、错误率等指标
  • 压力测试验证了高并发下的稳定性

在这里插入图片描述

八、总结

流式输出是 AI Agent 从"能用"到"好用"的关键一跃。它不仅仅是技术实现的问题,更是用户体验设计的核心环节。

回顾本文的核心内容:

第一,SSE 和 WebSocket 是两种主流的流式输出方案。 SSE 基于 HTTP,实现简单,适合服务器→客户端的单向流式场景;WebSocket 是全双工协议,适合需要双向实时通信的复杂场景。选择哪个方案,取决于你的 Agent 是否需要支持流式过程中的用户交互。对于大多数对话式 Agent,SSE 已经足够;对于需要中断、多 Agent 协作、实时通知的场景,WebSocket 更合适。

第二,前端渲染是体验的关键。 流式 Markdown 渲染需要处理不完整语法的边界情况——代码块未闭合、列表未结束、表格未完成等。使用 requestAnimationFrame 合并渲染请求是性能优化的关键手段。智能滚动检测确保用户向上查看历史时不会被自动滚动打断。

第三,中断处理是生产级的必需品。 一个完整的 Agent 中断流程包括:取消 LLM 生成任务、取消工具调用任务、执行资源清理回调、保留部分输出、通知前端。每个环节都需要超时保护,防止异常情况下的资源泄漏。AgentSession 数据结构提供了清晰的状态管理和资源追踪机制。

第四,技术方案的选择要回归业务需求。 不要为了用 WebSocket 而用 WebSocket。如果你的 Agent 只需要简单的对话流式输出,SSE 的实现成本更低、维护更简单。如果你的 Agent 需要复杂的交互(中断、追加指令、多步骤协作),WebSocket 的双向通信能力会带来更好的用户体验。

第五,生产环境需要系统性的保障措施。 连接管理(超时、心跳、重连)、安全防护(XSS 转义、鉴权、CSP)、性能优化(渲染合并、限流)、监控告警(连接数、延迟、错误率),这些缺一不可。本文提供的检查清单可以作为上线前的最后防线。

流式输出的本质是让 Agent 的"思考过程"变得可见,让用户从被动的等待者变成主动的参与者。当用户能实时看到 Agent 在做什么、能随时叫停、能即时调整方向,Agent 就从一个黑盒工具变成了一个透明的协作者。这种透明度,是用户信任 Agent 的基础。


参考资料

  1. MDN Web Docs - Server-Sent Events (SSE) — SSE 标准文档
  2. MDN Web Docs - WebSocket API — WebSocket 标准 API
  3. FastAPI StreamingResponse 官方文档 — FastAPI 流式响应
  4. Socket.IO 官方文档 — Socket.IO 完整指南
  5. OpenAI API - Streaming Responses — OpenAI 流式 API
  6. HTML Living Standard - Section 9.2 Server-sent events — SSE 规范
  7. RFC 6455 - The WebSocket Protocol — WebSocket 协议规范
  8. marked.js - Markdown Parser — 前端 Markdown 解析库
  9. Python-SocketIO 文档 — Python Socket.IO 实现
  10. nginx SSE 配置指南 — Nginx 代理缓冲配置

转载自 CSDN-专业IT技术社区

原文链接:https://blog.csdn.net/sinat_41617212/article/details/166903865

文章来源转载

评论

赞0

评论列表

微信小程序
QQ小程序

关于作者

点赞数:0
关注数:0
粉丝:0
文章:0
关注标签:0
加入于:--