AI编程 2026-07-31 83 次浏览

WebSocket 入门:用 Java 和 Python 跑通 Agent 实时事件流

从 WebSocket 双向通信原理出发,用 Java 21 与 Python 3.10+ 实现同一套可运行的 Agent 实时事件协议,覆盖流式回答、工具调用、任务取消与断线恢复设计。

摘要:当 Agent 不只需要流式回答,还要实时展示工具调用、接收任务取消和人工审批时,WebSocket 往往比单向推送更合适。本文从协议原理讲起,并用 Java 21 与 Python 3.10+ 跑通同一套可交叉验证的 Agent 事件协议。

WebSocket 在用户界面与 Agent 之间建立双向事件通道,并由 Java 与 Python 共同实现

聊天界面里,模型一边生成内容,页面一边出现文字,这只是 Agent 实时交互的第一步。

真正的 Agent 还可能搜索资料、执行代码、等待用户批准,甚至在执行途中被取消。此时,前端和 Agent 网关之间不再是单向的“服务端推送”,而是一条持续交换事件的通道。

这篇文章会完成三件事:理解 WebSocket 为什么适合这类场景;设计一套最小 Agent 事件协议;分别运行 Java 和 Python 两个版本,并用同一客户端交叉验证。

完整代码位于本文配套目录:examples/websocket-agent-demo/

完成后,你会得到什么

配套示例包含两个独立项目:

websocket-agent-demo/
├── java/       # Java 21 + Spring Boot 4.1
└── python/     # Python 3.10+ + websockets 16.1.1

两套服务实现完全相同的事件流程:

run.start
    ↓
run.status
    ↓
tool.call
    ↓
tool.result
    ↓
run.delta(多次)
    ↓
run.completed

客户端还可以在运行过程中发送 run.cancel,服务端会取消对应任务并返回 run.cancelled

这里暂时不接入具体大模型。示例使用一个可观察的模拟 Agent,把注意力集中在连接、并发和事件协议上。理解这条链路后,只需要把模拟生成函数替换为真实模型的流式调用。

WebSocket 是升级后的长期双向连接

WebSocket 的连接从 HTTP 请求开始。客户端带上 Upgrade: websocket 请求升级协议,服务端同意后返回 101 Switching Protocols

握手完成后,底层连接继续保持。客户端和服务端都可以主动发送文本帧或二进制帧,不必为每条消息重新建立 HTTP 请求。

客户端                              服务端
   │── HTTP Upgrade ───────────────▶│
   │◀── 101 Switching Protocols ────│
   │                                │
   │── run.start ──────────────────▶│
   │◀── tool.call ──────────────────│
   │◀── run.delta ──────────────────│
   │── run.cancel ─────────────────▶│

Spring 官方文档将 WebSocket 描述为建立在单个 TCP 连接上的全双工、双向通信通道,同时强调它只负责传输,不规定业务消息的语义。换句话说,run.starttool.callseq 都需要应用自己定义。参考:Spring WebSocket 文档

WebSocket 通过一次握手建立持续保持的双向连接

这张图最重要的结论不是“连接一直开着”,而是双方都能主动发事件。这正是任务取消、工具审批和实时语音块等场景需要的能力。

先判断方向,再决定是否使用 WebSocket

“需要流式输出”并不等于“必须使用 WebSocket”。

如果客户端只负责提交一次问题,之后只接收服务端输出,那么 HTTP 创建任务加服务器发送事件(Server-Sent Events,SSE)通常更简单。只有当执行期间存在频繁的双向事件时,WebSocket 的价值才会明显。

方案通信方向更适合的场景
HTTP请求—响应创建任务、查询历史、普通接口
SSE服务端到客户端Token 流、进度通知
WebSocket双向取消、审批、实时协作、音频块
消息队列服务端之间持久任务、重试、削峰、多 Agent 调度

Agent 通信应根据单向推送、双向交互和可靠后台任务选择方案

本文采用 WebSocket,是因为协议同时包含服务端主动推送和客户端主动控制。它不意味着 WebSocket 应该替代 REST 接口或消息队列。

一个常见的生产组合反而是:REST 创建和查询任务,WebSocket 负责界面实时交互,消息队列负责后台可靠执行。

用统一事件信封约束 Agent 消息

WebSocket 不会帮我们路由业务消息。因此,开始写处理器之前,先定义一个稳定的事件信封。

客户端启动任务:

{
  "version": 1,
  "type": "run.start",
  "runId": "run-001",
  "input": "介绍一下 WebSocket"
}

服务端增量返回:

{
  "version": 1,
  "type": "run.delta",
  "runId": "run-001",
  "seq": 4,
  "timestamp": "2026-07-29T03:17:05Z",
  "data": {
    "text": "建立了一条长期保持的"
  }
}

这些字段各有作用:

  • version:让协议以后可以演进。
  • type:决定消息由哪个处理逻辑消费。
  • runId:标识一次 Agent 运行,而不是一次网络连接。
  • seq:记录同一次运行内的事件顺序,为重连补发提供依据。
  • timestamp:便于排查跨服务的时间线。
  • data:承载不同事件自己的内容。

事件类型不要只设计一个模糊的 message。至少应该区分运行状态、工具调用、增量文本、结束、取消和错误。这样前端不需要从自然语言里猜测 Agent 正在做什么。

Java 版本:用 Spring Boot 注册原生处理器

Java 示例使用:

  • JDK 21
  • Maven 3.9+
  • Spring Boot 4.1.0
  • Spring Framework 原生 TextWebSocketHandler

Spring Boot 官方为 MVC WebSocket 提供了 spring-boot-starter-websocket。本文使用当前的原生处理器,没有额外引入 STOMP。参考:Spring Boot WebSocket 文档

进入 Java 目录并启动服务:

cd examples/websocket-agent-demo/java
mvn spring-boot:run

服务地址为:

ws://localhost:8080/ws/agent

注册端点的核心代码很短:

@Configuration
@EnableWebSocket
public class WebSocketConfig implements WebSocketConfigurer {

    private final AgentWebSocketHandler handler;

    public WebSocketConfig(AgentWebSocketHandler handler) {
        this.handler = handler;
    }

    @Override
    public void registerWebSocketHandlers(WebSocketHandlerRegistry registry) {
        registry.addHandler(handler, "/ws/agent");
    }
}

处理器收到 run.start 后,不直接在 WebSocket I/O 线程里执行耗时逻辑,而是提交给虚拟线程执行器:

private final ExecutorService workers =
        Executors.newVirtualThreadPerTaskExecutor();

@Override
protected void handleTextMessage(
        WebSocketSession session,
        TextMessage message
) {
    Command command = jsonMapper.readValue(
            message.getPayload(),
            Command.class
    );

    if ("run.start".equals(command.type())) {
        startRun(session, command);
    } else if ("run.cancel".equals(command.type())) {
        cancelRun(session, command.runId());
    }
}

完整代码还做了三件容易被最小示例忽略的事:

  1. sessionId + runId 定位正在运行的任务。
  2. 收到取消事件时调用 FutureTask.cancel(true)
  3. 对同一个 WebSocketSession 的发送操作进行串行化。

第三点很重要。Agent 可能同时产生模型增量、工具日志和状态事件,不能假设多个线程可以随意向同一个会话并发写入。本文为了突出原理使用 synchronized;正式项目可以进一步使用单会话发送队列、背压控制或 Spring 的并发会话装饰器。

另一个版本差异是 JSON 库。Spring Boot 4 默认使用 Jackson 3,因此完整代码注入的是:

import tools.jackson.databind.json.JsonMapper;

如果项目仍是 Spring Boot 3,常见类型是 com.fasterxml.jackson.databind.ObjectMapper。Spring Boot 4 的 JSON 迁移说明可参考:Spring Boot JSON 文档

另开终端运行 Java 客户端:

cd examples/websocket-agent-demo/java
mvn -q compile exec:java \
  -Dexec.mainClass=demo.ws.AgentClient \
  -Dexec.args="请介绍 WebSocket"

Java 客户端使用 JDK 自带的 java.net.http.WebSocket。该客户端 API 从 Java 11 开始提供,并通过 HttpClient.newWebSocketBuilder() 异步建立连接。参考:Oracle WebSocket.Builder API

Python 版本:用 asyncio 管理连接和运行任务

Python 示例使用:

  • Python 3.10+
  • websockets==16.1.1
  • asyncio.Task 管理每次 Agent 运行

创建虚拟环境并安装依赖:

cd examples/websocket-agent-demo/python
python3 -m venv .venv
source .venv/bin/activate
python -m pip install -r requirements.txt
python server.py

服务地址为:

ws://localhost:8765/ws/agent

Python 服务为每个连接维护一个 AgentConnection,为每个 runId 保存一个任务:

async def start_run(self, run_id, input_text):
    run_id = run_id or str(uuid.uuid4())

    task = asyncio.create_task(
        self.simulate_agent(run_id, input_text)
    )

    self.tasks[run_id] = task

收到取消事件时调用:

task = self.tasks.get(run_id)
task.cancel()

模拟 Agent 捕获 asyncio.CancelledError,先发送 run.cancelled,再结束任务。真实模型调用也应该把取消信号继续传递到底层请求,否则界面看似取消了,后端模型和工具仍可能继续消耗资源。

另开终端运行客户端:

cd examples/websocket-agent-demo/python
source .venv/bin/activate
python client.py "请介绍 WebSocket"

当前 websockets 官方文档推荐从 websockets.asyncio.server 导入 serve,客户端则可以使用 connect() 异步上下文管理器。参考:websockets 运行示例

用跨语言调用验证协议没有绑死实现

Java 和 Python 服务的端口不同,但事件结构一致。因此,Python 客户端可以直接访问 Java 服务:

python client.py \
  "Python 客户端访问 Java 服务" \
  ws://localhost:8080/ws/agent

正常情况下会依次看到:

run.status       seq=1
tool.call        seq=2
tool.result      seq=3
run.delta        seq=4
run.delta        seq=5
run.delta        seq=6
run.completed    seq=7

这一步验证的不是语言性能,而是协议边界:客户端只依赖事件格式,不依赖服务端究竟由 Java 还是 Python 实现。

Java 与 Python 服务使用同一事件协议,客户端可以交叉访问

配套代码已经完成以下实际验证:Java 项目编译并启动成功;Java 客户端访问 Java 服务成功;Python 客户端访问 Python 服务成功;Python 客户端访问 Java 服务成功。

接入真实 Agent 时,只替换事件生产者

示例中的 simulateAgent 会依次发送工具事件和几个固定文本片段。接入真实 Agent 时,连接层和事件信封可以保持不变,只替换事件来源:

模型输出 Token        → run.delta
Agent 开始调用工具    → tool.call
工具返回结果          → tool.result
等待用户批准          → tool.approval_required
正常结束              → run.completed
用户取消              → run.cancelled
异常结束              → error

用户批准工具调用时,可以从客户端发送:

{
  "version": 1,
  "type": "tool.approve",
  "runId": "run-001",
  "callId": "call-8"
}

这也是 WebSocket 相比纯 SSE 更有价值的地方:审批、取消和补充输入不必另开一套实时控制通道。

WebSocket 连接不等于 Agent 运行

这是从 Demo 走向生产最重要的设计边界。

手机切换网络、浏览器休眠、负载均衡器回收空闲连接,都可能让 WebSocket 断开。但 Agent 任务可能仍在运行。如果把所有状态都存放在 WebSocketSession 或 Python 连接对象里,连接一断,任务历史和恢复位置也会一起丢失。

WebSocket 可以断开重连,而 Agent 运行需要依靠运行编号和事件序号持续存在

生产系统应该把两者分开:

  • connectionId 表示一条可能随时断开的网络连接。
  • runId 表示一次需要持续追踪的 Agent 运行。
  • seq 表示该运行已经产生到第几个事件。
  • 事件历史保存在数据库、Redis Streams 或消息系统中。

重连后,客户端可以携带 runId 和最后收到的 seq,服务端先补发缺失事件,再继续实时推送。

上生产前补齐六个能力

1. 认证和来源校验

生产环境使用 wss://。握手阶段验证用户身份和 Origin,不要为了省事开放所有来源。

浏览器原生 WebSocket API 不方便设置任意 Authorization 请求头,可以使用安全 Cookie 或短期、一次性的连接凭证。不要把长期令牌直接放在 URL 中,因为 URL 可能进入日志和监控系统。

2. 心跳和空闲超时

使用 Ping/Pong 或应用层心跳发现失效连接,同时配置合理的空闲超时。只依赖 TCP 感知断线,反馈可能不够及时。

3. 消息大小和背压

限制单条消息大小和待发送队列。慢客户端如果长期跟不上 Token 和日志产生速度,服务端必须选择暂停、合并、丢弃非关键事件或关闭连接。

4. 断线重连和事件补发

客户端采用带抖动的指数退避重连。服务端使用 runId + seq 恢复,而不是把重连误认为一次新任务。

5. 多实例状态共享

多实例部署时,不要只把运行状态放在进程内存里。可以使用 Redis、数据库或消息队列共享运行状态,并确认网关、反向代理和负载均衡器支持 WebSocket Upgrade。

6. 可观察性

日志至少关联 userIdconnectionIdrunIdeventIdseq。这样才能判断问题发生在模型、工具、事件总线还是网络发送阶段。

三个常见误区

只要流式输出就使用 WebSocket

如果只有服务端向客户端发送 Token,优先评估 SSE。它的浏览器接入、代理兼容和重连语义通常更直接。

用 WebSocket 代替后台消息队列

WebSocket 的目标是实时通信,不天然提供持久化、消费确认和失败重试。多 Agent 后台调度仍应使用适合的任务系统或消息队列。

一开始就引入 STOMP

STOMP 可以提供订阅、主题和消息代理语义,适合聊天室、广播和复杂订阅。点对点的 Agent 事件流可以先使用原生 WebSocket,把业务协议理解清楚后再决定是否增加协议层。

总结:把 WebSocket 放在正确的位置

WebSocket 最适合承担的是用户界面与 Agent 网关之间的实时双向事件通道

记住四个结论:

  1. 只有双向实时控制明显时,WebSocket 才比 SSE 更有优势。
  2. WebSocket 只负责传输,Agent 必须定义稳定的事件信封和事件类型。
  3. Java 与 Python 可以共享同一协议,客户端不应该依赖服务端语言。
  4. 连接可能断开,Agent 运行必须通过 runId + seq 独立保存和恢复。

下一步可以把示例中的 simulateAgent 替换为真实模型流,并优先补上 tool.approval_requiredrun.cancel 和断线续传。完成这三项后,它才开始接近一条真正可用的 Agent 实时链路。