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

聊天界面里,模型一边生成内容,页面一边出现文字,这只是 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.start、tool.call 和 seq 都需要应用自己定义。参考:Spring WebSocket 文档。

这张图最重要的结论不是“连接一直开着”,而是双方都能主动发事件。这正是任务取消、工具审批和实时语音块等场景需要的能力。
先判断方向,再决定是否使用 WebSocket
“需要流式输出”并不等于“必须使用 WebSocket”。
如果客户端只负责提交一次问题,之后只接收服务端输出,那么 HTTP 创建任务加服务器发送事件(Server-Sent Events,SSE)通常更简单。只有当执行期间存在频繁的双向事件时,WebSocket 的价值才会明显。
| 方案 | 通信方向 | 更适合的场景 |
|---|---|---|
| HTTP | 请求—响应 | 创建任务、查询历史、普通接口 |
| SSE | 服务端到客户端 | Token 流、进度通知 |
| WebSocket | 双向 | 取消、审批、实时协作、音频块 |
| 消息队列 | 服务端之间 | 持久任务、重试、削峰、多 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());
}
}
完整代码还做了三件容易被最小示例忽略的事:
- 用
sessionId + runId定位正在运行的任务。 - 收到取消事件时调用
FutureTask.cancel(true)。 - 对同一个
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.1asyncio.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 项目编译并启动成功;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 连接对象里,连接一断,任务历史和恢复位置也会一起丢失。

生产系统应该把两者分开:
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. 可观察性
日志至少关联 userId、connectionId、runId、eventId 和 seq。这样才能判断问题发生在模型、工具、事件总线还是网络发送阶段。
三个常见误区
只要流式输出就使用 WebSocket
如果只有服务端向客户端发送 Token,优先评估 SSE。它的浏览器接入、代理兼容和重连语义通常更直接。
用 WebSocket 代替后台消息队列
WebSocket 的目标是实时通信,不天然提供持久化、消费确认和失败重试。多 Agent 后台调度仍应使用适合的任务系统或消息队列。
一开始就引入 STOMP
STOMP 可以提供订阅、主题和消息代理语义,适合聊天室、广播和复杂订阅。点对点的 Agent 事件流可以先使用原生 WebSocket,把业务协议理解清楚后再决定是否增加协议层。
总结:把 WebSocket 放在正确的位置
WebSocket 最适合承担的是用户界面与 Agent 网关之间的实时双向事件通道。
记住四个结论:
- 只有双向实时控制明显时,WebSocket 才比 SSE 更有优势。
- WebSocket 只负责传输,Agent 必须定义稳定的事件信封和事件类型。
- Java 与 Python 可以共享同一协议,客户端不应该依赖服务端语言。
- 连接可能断开,Agent 运行必须通过
runId + seq独立保存和恢复。
下一步可以把示例中的 simulateAgent 替换为真实模型流,并优先补上 tool.approval_required、run.cancel 和断线续传。完成这三项后,它才开始接近一条真正可用的 Agent 实时链路。