这周我们把一个上线后被吐槽最多的 Agent 模块彻底重写了一遍。用户的原话是:“模型思考都开始了,为什么前端还在疯狂转圈?”你仔细去看就会发现问题不在大模型,而在传输链路——Agent 一边流式输出思考过程,一边去调工具,这个过程是典型的“长时、低频、单向下推”数据流,用传统轮询去做,短了浪费服务器资源,长了体验像断流。这篇文章是我在 Spring Boot 3 + SSE + Redis 上的一套落地解法,完整覆盖流式思考、工具调用事件、跨实例分发和断线重连,适合正在做 Agent 应用、或者想把手里的轮询接口改造成实时推送的朋友。
1. 为什么是SSE而不是WebSocket:流式Agent场景的技术选型复盘
很多团队一提到实时推送,第一反应就是 WebSocket。其实 Agent 这个场景跟聊天室不一样,它没有强双向通信需求,服务端把思考片段、工具状态、最终结果一条条推给客户端就够了。SSE(Server-Sent Events)是跑在 HTTP 上的单向流式协议,对服务端压力比 WebSocket 小很多,还天然支持自动重连,这三点刚好对上 Agent 输出的脾气。
1.1 轮询与“假思考”:真实体验问题的拆解
我见过不少 Agent 演示站,前端是每 2 秒轮询一次任务状态接口。这个方案在任务 5 秒内跑完的时候勉强能看,但一旦你的 Agent 接入多个工具,一个任务动辄二十秒、甚至半分钟以上,用户的体验就是:点一下按钮,转圈,等 2 秒刷新,发现还在“思考中”,再过 4 秒,突然一下蹦出结果。中间过程完全是黑盒。
更难受的是轮询节奏。你设 2 秒一次,工具调用卡了 5 秒,用户就盯着页面发呆;你设 0.5 秒一次,Nginx、网关、应用线程全在空转,高峰时一个一千并发量的 Agent 应用,轮询请求就能打满服务端的线程池。我线上遇到过一台 4C8G 的机器,active connections 没过两百,但 Tomcat 线程池已经烧到 90%,原因就是一堆客户端在无脑轮询。
轮询还掩盖了一个问题:你没法判断任务到底是在思考、在调工具,还是已经死了。接口返回永远是pending,用户只会觉得“这 AI 是不是卡了”。说白了,轮询解决的是“我定期去看一眼结果”,但 Agent 场景需要的是“你有了进展就主动告诉我”。
1.2 SSE、WebSocket与轮询的取舍对比
我们拿一个标准 Agent 任务来说:用户问一个需要搜索+计算+推理的问题,LLM 先输出思考过程,然后发起工具调用,工具执行完成后 LLM 继续推理,最后输出结论。这个过程中,客户端要的是服务端单向的多个事件,几乎没有上行数据需求。
SSE 在这种场景下有四个实打实的优点:
- HTTP 协议自带,不用额外握手和升级,防火墙、网关、负载均衡对它的兼容度比对 WebSocket 高得多。
- 自动重连是协议层行为。WebSocket 断线了得自己写心跳和重连逻辑;SSE 的
EventSource客户端碰到断连会自动重试,服务端只需要在重连时把缺失事件补回去。 - 服务端资源占用更低,WebSocket 需要保活一个双向长连接,SSE 底层还是普通 HTTP 响应,只是把 Content-Type 设为
text/event-stream,对 Tomcat 这类容器来说就是个普通异步响应。 - 开发心智简单。一个 GET 接口,返回
SseEmitter,事件往里面send()就行。
WebSocket 不是不好,但你得问自己一个问题:我需要客户端往服务端推数据吗?Agent 场景里的工具调用参数、取消指令,这些用普通 POST 请求完全够了。强行上 WebSocket,带来的反而是连接状态管理、心跳保活、消息格式自定义这些额外工作量。
1.3 Spring Boot 3 对SSE的支持现状
Spring Boot 3.x 里,支持 SSE 的路子主要有两条:Spring MVC 的SseEmitter,和 Spring WebFlux 的Flux<ServerSentEvent>。如果你们的工程本来就是传统的 Spring Boot 3 + Tomcat + 同步 JDBC 这套技术栈,SseEmitter是最平滑的接入方式,它基于 Servlet 3.1 异步响应,不需要把整个服务改造成 WebFlux 响应式模式。
WebFlux 的方案更适合一个从零开始、且全链路非阻塞的项目。它能让你用Flux<ServerSentEvent>直接返回一个流,配合 Reactor 的背压机制,代码会很简洁。但有一个容易踩的坑:如果你在 WebFlux 代码里混用了阻塞操作(比如 JDBC 同步查询、阻塞 Redis 客户端),线程模型会被破坏,性能反而不如 MVC。我的建议是:沿用你现有的技术栈,不要为了 SSE 专门换一个全家桶。我们后面实现主链路用的是SseEmitter这套。
2. 整体链路:Agent思考流、工具调用与Redis的分工
先把架构捋清楚。一个典型的 Agent 任务包含三件事:LLM 流式思考、工具调用、结果回填。这三件事产生的事件类型不同,但它们都得在一个 SSE 连接上、按顺序推到客户端。Redis 在这个链路里不是主角,但少了它,跨实例的水平扩展和任务状态管理就会非常痛苦。
2.1 把Agent输出拆成三类事件
我建议把所有输出统一成一个AgentEvent,用type字段区分:
thinking:LLM 推理过程中的 token 增量。用户能实时看到 AI 在“想什么”。tool_start/tool_result:Agent 准备调用某个工具、以及工具返回结果。这是“工具调用中枢”的关键,用户能看到 AI 在查什么、查到了什么。done/error:任务结束或异常退出。
事件里还要带上taskId、seq序号、timestamp。seq特别重要,它是断线重连时对账的依据,后文会专门讲。
有人会把thinking事件每生成一个 token 就发一次,我实测下来这样做太激进。LLM 接口返回的是增量流,如果你直接把每个 token 都转成 SSE 消息,Redis 的 Pub/Sub 吞吐会暴涨,前端渲染也跟不上。比较稳妥的做法是小批量聚合,比如累积 30~50 毫秒的 token,或者按障碍符号(换行、句号)切割成片段再发。看板上的效果跟逐字差不多,但服务端压力小一个数量级。
2.2 Redis的三个角色:事件通道、任务状态、结果缓冲
这条链路里 Redis 同时干了三件事:
- 事件通道:所有 Agent 执行事件通过 Redis Pub/Sub 广播出去,持有对应 SSE 连接的实例收到后推给客户端。这叫“跨实例分发”。
- 任务状态:把一个任务的完整生命周期存成 Hash,状态机流转都写在这里,前端重建页面或服务重启后能查到任务去哪了。
- 结果缓冲:工具调用的返回结果、Agent 最终回答,先写到 Redis 再异步落库,避免高频事件直接打到数据库。
有些人会问:为什么不直接用消息队列,比如 RabbitMQ 或者 Kafka?如果你们公司 MQ 基建已经很成熟,用 MQ 当然更稳。但 Redis Pub/Sub 的优势是零新组件,Spring Boot 项目里反正都要用 Redis,链路越短越好排查。当然,Pub/Sub 消息不落地、没有持久化,所以生产环境我会把事件同时XADD到 Redis Stream 里做最近 30 分钟的历史重放,这部分在第 5 节展开。
2.3 一次完整任务的生命周期
用一个表格描述典型链路:
| 阶段 | 动作 | 数据流 |
|---|---|---|
| 1 | 客户端 POST /api/agent/tasks | 生成 taskId,任务状态置为 pending,返回给前端 |
| 2 | 客户端 GET /stream 建立 SSE | 服务端注册 SseEmitter,订阅 Redis 全局事件通道 |
| 3 | Agent 执行器开始跑 | LLM token 聚合成 thinking 事件,发布到 Redis |
| 4 | 触发工具调用 | 发 tool_start,执行工具,发 tool_result,继续调 LLM |
| 5 | 推理完成 | 发 done 事件,更新 Redis 任务状态为 succeeded |
| 6 | 客户端收齐 | 关闭 SSE,按需落库 |
这里要注意,第 1 步和第 2 步之间是有竞态的:任务可能在 SSE 建立之前就已经跑完或者报错了。所以打开流的时候,服务端第一步不能直接开启实时订阅,得先去 Redis 里读一下任务状态,如果已经是终态,直接把 done 或者 error 的事件补发,然后再订阅实时通道。这个细节不做,用户偶尔会看到任务“凭空消失”。
2.4 多实例横向扩展的形态
SSE 长连接是有状态的,如果部署了多个应用实例,客户端连上实例 A,而 Agent 任务可能在实例 B 上执行。Redis Pub/Sub 在这里就成了粘合剂:任务执行器把事件发布到全局通道agent:event:global,所有实例都订阅这个通道,只有本地注册表里存在该 taskId 对应 emitter 的实例才真正把事件推到客户端。
这套方案在十来个节点以内都够用。事件总量大、订阅实例多导致广播放大时,再考虑把全局通道拆成agent:event:{taskId}按任务路由的版本。
3. 关键实现一:Spring Boot 3 + SSE 服务端与客户端接入
下面进入代码,看完可以直接往你工程里搬。项目基于 Spring Boot 3.3、Java 17、Redis 7。
3.1 基于SseEmitter的流式端点
先写 Controller,核心是 POST 创建任务和 GET 拉流两个接口:
@RestController @RequestMapping("/api/agent") public class AgentTaskController { private final AgentTaskService taskService; public AgentTaskController(AgentTaskService taskService) { this.taskService = taskService; } @PostMapping("/tasks") public ResponseEntity<TaskCreateResult> createTask(@RequestBody AgentTaskRequest request) { String taskId = taskService.submit(request); return ResponseEntity.accepted().body(new TaskCreateResult(taskId)); } @GetMapping(value = "/tasks/{taskId}/stream", produces = MediaType.TEXT_EVENT_STREAM_VALUE) public SseEmitter stream(@PathVariable String taskId, @RequestHeader(value = "Last-Event-ID", required = false) Long lastSeq) { return taskService.openStream(taskId, lastSeq); } @PostMapping("/tasks/{taskId}/cancel") public ResponseEntity<Void> cancel(@PathVariable String taskId) { taskService.cancel(taskId); return ResponseEntity.noContent().build(); } }SseEmitter构造参数传0L表示永不超时。这里容易出问题的是:一旦服务端自己不超时,如果客户端断开没人处理,连接会一直挂着。所以必须显式注册onCompletion、onTimeout、onError三个回调,把 emitter 从本地注册表清掉。我见过不少人只注册了 onCompletion 没管 onError,结果客户端断电后本地注册表里全是僵尸连接。
3.2 事件结构与协议设计
事件体设计成 record,直接走 JSON:
public record AgentEvent( Long seq, String taskId, String type, // thinking / tool_start / tool_result / done / error / ping String content, Map<String, Object> meta, long timestamp ) {}meta里面放结构化的附加信息。比如工具调用事件,meta里装toolName、arguments、durationMs;done事件里放usage、model。前端拿到后,既能展示纯文本,也能做一些结构化渲染。
发送事件时,SSE 协议允许带id字段。我们把seq放在事件 id 上,这样浏览器原生EventSource断线重连时,会自动在请求头带上Last-Event-ID,服务端就能据此重放缺失事件。代码:
private SseEmitter.SseEventBuilder toSseEvent(AgentEvent event) { return SseEmitter.event() .id(String.valueOf(event.seq())) .name(event.type()) .data(JsonUtils.toJson(event)); }这里用.name()对应 SSE 的event:字段,前端就能用addEventListener('thinking', ...)分别监听不同类型,而不需要自己解析type。
3.3 多实例场景下用Redis Pub/Sub交叉分发
关键是“发布”和“订阅”两件事要配套。
发布:
@Service public class AgentEventPublisher { private static final String GLOBAL_CHANNEL = "agent:event:global"; private final RedisTemplate<String, String> redisTemplate; public void publish(AgentEvent event) { // 统一过 Redis 通道,保证任意实例上的执行器都能把事件送给持有 SSE 连接的实例 redisTemplate.convertAndSend(GLOBAL_CHANNEL, JsonUtils.toJson(event)); } }订阅用 Boot 自带的RedisMessageListenerContainer:
@Configuration public class RedisPubSubConfig { @Bean public RedisMessageListenerContainer redisMessageListenerContainer( RedisConnectionFactory connectionFactory, AgentEventSubscriber subscriber) { RedisMessageListenerContainer container = new RedisMessageListenerContainer(); container.setConnectionFactory(connectionFactory); container.addMessageListener(subscriber, new PatternTopic("agent:event:global")); // 这个线程池大小决定事件处理上限,后面第5节会展开说 container.setTaskExecutor(Executors.newFixedThreadPool(Runtime.getRuntime().availableProcessors() * 2)); return container; } }订阅处理:
@Component public class AgentEventSubscriber implements MessageListener { private final LocalEmitterRegistry localRegistry; @Override public void onMessage(Message message, byte[] pattern) { AgentEvent event = JsonUtils.fromJson(message.getBody(), AgentEvent.class); List<SseEmitter> emitters = localRegistry.get(event.taskId()); if (emitters == null || emitters.isEmpty()) { return; } for (SseEmitter emitter : emitters) { sendToEmitter(emitter, event); } } }LocalEmitterRegistry是一个ConcurrentHashMap<String, CopyOnWriteArrayList<SseEmitter>>,key 是 taskId。要特别注意:同一个 taskId 可能对应多个 SSE 连接,因为用户可能开两个标签页,也可能重连接入。如果注册表只存单个 emitter,第二个标签页就收不到事件。
3.4 前端EventSource接入与自动重连
前端不用装什么重库,原生EventSource就行:
const es = new EventSource(`/api/agent/tasks/${taskId}/stream`); es.addEventListener('thinking', (e) => { appendThinking(JSON.parse(e.data)); }); es.addEventListener('tool_start', (e) => { appendTool(JSON.parse(e.data)); }); es.addEventListener('tool_result', (e) => { appendToolResult(JSON.parse(e.data)); }); es.addEventListener('done', (e) => { renderFinal(JSON.parse(e.data)); es.close(); }); es.addEventListener('error', (e) => { // EventSource 断线会自动重连,不用在这里手动 new // 但超过一定次数或业务上决定放弃时,需要主动 close });原生 EventSource 断线重连是自动的。但如果中途代理层把连接掐了、且重连请求里带了Last-Event-ID,后端应能从历史事件补齐中间缺口。这里有个细节:关闭连接一定要调es.close(),不然浏览器会按它的重试策略继续打你的后端。我遇到过一个线上案例,用户关掉页面但 WebSocket 连接还活着,服务端状态机迟迟不清理——SSE 是单向的,浏览器页面关了,TCP 并不会立刻感知到,服务端要自己依靠心跳或回调判断。
4. 关键实现二:Redis状态管理与工具调用闭环
SSE 解决的是“怎么把事件送出去”,这一节解决“Agent 执行过程中,状态和工具调用怎么跟事件流咬合”。
4.1 任务状态机与Key设计
我用一个 Hash 存任务状态:
| Key | 类型 | 说明 | TTL |
|---|---|---|---|
agent:task:{taskId} | Hash | status、createdAt、model、errorMsg、usage | 30 分钟 |
agent:event:{taskId} | Stream | 事件历史,用于断线重放 | 30 分钟 |
agent:lock:{taskId} | String | 分布式锁,防止重复提交 | 创建任务时设置 10s |
Hash 里的status字段取值很简单:pending -> running -> succeeded | failed | canceled。
工具调用结果这类临时数据,我一般放在agent:tool:{taskId}:{callId},TTL 设 10 分钟。为什么要单独存?因为工具返回结果要回填给 LLM 上下文,如果执行器和后续推理不在一台机器上,需要从 Redis 取。Agent 跑完后这些数据也就没用了,TTL 到期自动清理,不占地方。
4.2 工具调用的发起与结果回填
核心执行循环大致长这样:
public void execute(String taskId, AgentTaskRequest request) { // 更新状态为 running stateMachine.update(taskId, "running"); List<ChatMessage> messages = new ArrayList<>(request.messages()); while (true) { // 流式调用大模型,按聚合策略发 thinking 事件 StreamResult stream = llmService.chatStream(messages); stream.tokensChunks().forEach(chunk -> { publisher.publish(new AgentEvent(nextSeq(taskId), taskId, "thinking", chunk.text(), buildMeta(stream), System.currentTimeMillis())); }); // 判断模型响应里是否包含工具调用 if (stream.hasToolCalls()) { for (ToolCall call : stream.toolCalls()) { publisher.publish(new AgentEvent(nextSeq(taskId), taskId, "tool_start", call.name(), buildToolMeta(call), System.currentTimeMillis())); // 执行工具的阻塞调用 Object result = toolExecutor.execute(call.name(), call.arguments()); // 工具结果写 Redis,也发一个 tool_result 事件 redisTemplate.opsForValue().set( String.format("agent:tool:%s:%s", taskId, call.id()), JsonUtils.toJson(result), Duration.ofMinutes(10)); publisher.publish(new AgentEvent(nextSeq(taskId), taskId, "tool_result", call.name(), buildToolResultMeta(call, result), System.currentTimeMillis())); // 把工具结果拼进上下文,让模型基于工具输出做下一步推理 messages.add(new ToolMessage(call.id(), JsonUtils.toJson(result))); } continue; // 进入下一轮推理 } // 没有工具调用,说明推理已完成 String finalAnswer = stream.finalText(); publisher.publish(new AgentEvent(nextSeq(taskId), taskId, "done", finalAnswer, buildFinalMeta(stream), System.currentTimeMillis())); stateMachine.update(taskId, "succeeded"); break; } }这是目前 Agent 框架(LangChain、Dify、CrewAI 这类)里最常见的 ReAct 循环变体。不管你是直接用官方 SDK 还是自己封装,思路都是一样的:LLM 输出要么是普通文本,要么是工具调用指令,你把这个循环的每一个关键步骤都切成事件发出去,SSE 中枢就活了。
工具调用这个地方要多说一句:绝大多数工具调用是同步阻塞的,比如查数据库、调第三方 API。如果直接放在主线程跑,前面所有 thinking 事件都是好的,但到了工具这一步就卡住,用户会看到“思考到一半停住了”。所以执行器最好开一个独立的线程池跑 Agent 循环,SSE 发送线程和 Agent 执行线程分离。SseEmitter.send()本身是线程安全的,可以从任意线程调。
4.3 幂等控制与重复提交防护
用户手一抖点了三次提交,后端就可能创建三个任务。这个我用 Redis 分布式锁挡住:
public String submit(AgentTaskRequest request) { String lockKey = "agent:lock:" + request.clientTraceId(); Boolean locked = redisTemplate.opsForValue() .setIfAbsent(lockKey, "1", Duration.ofSeconds(10)); if (!Boolean.TRUE.equals(locked)) { throw new BizException("任务正在处理中,请勿重复提交"); } String taskId = idGenerator.nextId(); saveTaskState(taskId, "pending"); executeAsync(taskId, request); return taskId; }锁的作用不是为了并发安全,而是为了用户体验。clientTraceId是前端生成的唯一标识,同一用户同一意图的重复点击会被锁拦下来。锁的 TTL 设 10 秒,是因为创建任务本身就很快,如果 10 秒还没处理完,说明链路后端有问题,这时候再放行几个重复任务反而是合理的。
工具回填那边也有一个幂等点:有些工具是异步回调的,比如发起一个耗时任务,结果通过 Webhook 回来。Webhook 可能重试,所以回填时要用SETNX或者判断任务状态,避免同一个callId的结果被写两次,导致 LLM 上下文里出现重复工具信息。
4.4 生产配置:序列化、TTL与连接池
这节全是血泪教训。
第一,序列化。默认的JdkSerializationRedisSerializer会让 Redis 里出现一堆\xAC\xED乱码,排查工具结果没法看,而且体积膨胀严重。我统一对事件和状态用StringRedisTemplate的 JSON 序列化,或者自定义GenericJackson2JsonRedisSerializer。调试的时候用redis-cli看数据,一眼就能明白。
第二,TTL。任务 Hash、事件 Stream 都必须设置 TTL。我碰到过一个事故,任务事件没设过期,Redis 内存被一堆已完成任务的 thinking 事件堆满,晚高峰触发了内存淘汰策略,导致业务 key 被误删。现在我的规则很固定:终态任务 TTL 30 分钟,中间状态的锁 TTL 5~10 秒,工具结果 TTL 10 分钟。
第三,连接池。Lettuce 默认配置在低并发下没问题,但 SSE 场景里动用 Redis 的频率很高,建议显式配置:
spring: data: redis: lettuce: pool: max-active: 32 max-idle: 16 min-idle: 4 max-wait: 200ms特别提醒:不要在主线程里同步调用 Redis 命令,尤其不要在SseEmitter的发送回调里写阻塞式 Redis 操作。一旦 Redis 抖动,整个 SSE 发送线程池都会被拖住,用户那边就是整片断流。要么用异步命令,要么把 Redis 调用放到独立线程池,二者必选其一。
5. 生产环境踩坑实录:心跳、背压与断线重连
这节的内容是我在真实线上环境一点一点踩出来的,也是“生产级”这三个字最值钱的部分。
5.1 Nginx缓冲导致的心跳假死
SSE 有一个非常经典的坑:Nginx 默认开着proxy_buffering,会把后端响应缓冲起来,攒到一定量才发给客户端。结果就是 LLM 已经出了一堆 token,用户那边几十秒没动静,然后突然一次性蹦出来,体验比轮询还差。
解决方法是双重保障。Nginx 配置:
location /api/agent/ { proxy_pass http://backend; proxy_buffering off; proxy_cache off; proxy_connect_timeout 5s; proxy_read_timeout 3600s; proxy_send_timeout 3600s; chunked_transfer_encoding off; }后端响应头也要带上:
@GetMapping(value = "/tasks/{taskId}/stream", produces = MediaType.TEXT_EVENT_STREAM_VALUE) public SseEmitter stream(@PathVariable String taskId) { HttpServletResponse response = ...; response.setHeader("X-Accel-Buffering", "no"); return service.openStream(taskId, null); }X-Accel-Buffering: no是 Nginx 官方提供的响应头开关,告诉 Nginx 别缓冲这个响应。网关后面还挂了云上负载均衡的话,需要在负载均衡那边也把缓冲关掉,这个各云厂商配置不一样,上线前一定要测首包延迟。
5.2 线程池耗尽与事件堆积
RedisMessageListenerContainer默认使用一个单线程处理所有订阅消息。我一度以为 Redis Pub/Sub 消息量小没问题,结果压测时发现事件处理线程成了瓶颈:Agent 并发二三十个任务,每个任务都是 thinking + 工具事件往外发,一个 Redis 订阅线程根本忙不过来,消息堆积导致前端事件延迟从几十毫秒涨到几秒。
解决就是给监听容器配置专门的线程池,数量不要拍脑袋。我的经验公式是CPU 核数 × 2,一般够用。事件处理线程只做“反序列化 + 注册表路由 + emitter.send()”,绝不在这里做业务逻辑。另外,如果某个 emitter 的send()抛出 IOException,说明客户端已经断开,要把这个 emitter 从注册表移除,不能让它反复报错拖慢线程。
5.3 断线重连后的事件重放
有一类问题很隐蔽:用户网络抖了一下,前端EventSource自动重连了,但重连期间 Agent 已经发了几条事件。由于 Redis Pub/Sub 不持久化,这几条事件永远丢了。用户看到的画面就是思路突然跳了一段,前面说“准备查天气”,下一个事件直接就是最终答案。
我最终用 Redis Stream 做事件历史。执行器每发一个事件,除了PUBLISH到全局通道,同时XADD写入agent:event:{taskId},容量控制用MAXLEN ~ 5000,过期时间靠EXPIRE30 分钟。客户端重连时带了Last-Event-ID,服务端就先XREAD读出该 seq 之后的事件全部补发,然后再订阅实时通道:
public SseEmitter openStream(String taskId, Long lastSeq) { SseEmitter emitter = new SseEmitter(0L); // 補发历史事件 if (lastSeq != null) { List<AgentEvent> history = eventHistory.rangeFrom(taskId, lastSeq); for (AgentEvent event : history) { safeSend(emitter, toSseEvent(event)); } } localRegistry.register(taskId, emitter); return emitter; }这个方案有个天然的竞态:补发历史和订阅实时之间可能有空隙,事件会重或者漏。要彻底解决得做 sequence 去重,前端收到事件后按seq去重即可,简单有效。我现在的前端实现就是维护一个lastSeq,只会接受seq > lastSeq的事件,重复的直接忽略。
5.4 如何定位一段8秒的“卡顿”
用户反馈“AI 思考卡了 8 秒”。这种问题排查起来其实不复杂,前提是你的事件里有足够多的上下文。我通常按这个顺序定位:
- 看前端你自己埋的事件时间线,卡顿发生在一个
thinking和下一个thinking之间,还是tool_start和tool_result之间。 - 如果是模型推理阶段卡顿,去看 LLM 请求日志,拿时间戳对比
thinking事件的时间戳,多半是模型首包延迟或者网络问题。 - 如果是工具阶段,用
tool_result里的durationMs字段,直接暴露是第三方 API 慢,还是 Redis 读写慢。 - 如果事件时间戳正常,但前端收到晚了,那就是传输链路问题,查 Nginx 缓冲、查订阅线程队列积压。
给事件打上timestamp这个习惯,能让你在出事的时候少掉一半头发。现在的 Agent 中间件和框架在可观测性上普遍弱,你自己不做事件打点,等线上出问题只能抓瞎。
6. 实测效果与可扩展性
一套方案好不好,最终要看数据和落地感受。
6.1 轮询 vs SSE的首包延迟对比
我在测试环境模拟了一个包含两次工具调用的 Agent 任务,总时长约 20 秒。轮询方案设置 2 秒间隔,用户从点击按钮到看到第一条思考内容,平均需要 2~6 秒;换成 SSE 后,首包延迟约 300~800 毫秒,主要消耗在任务创建和 SSE 连接建立上,一旦连接建好,思考事件是源源不断往下推的。
这个差异会在长任务上被指数级放大。工具调用多、推理轮次多的 Agent,总时长可能拉到 1 分钟以上,轮询的“假思考”体验会让人完全失去耐心,而 SSE 至少让用户看到 AI 正在一步一步推进。
6.2 并发与长连接压测数据
压测环境是两台 4C8G 的云主机,Redis 单独一台 2C4G。用脚本模拟 200 个并发 SSE 连接,每个连接维持 60 秒,服务端持续每隔 2 秒发一条事件。
结果:后端 CPU 稳定在 30% 左右,Tomcat 线程池利用率比轮询方案低了约 70%,Redis 的 Pub/Sub 吞吐在每秒 1000 条事件左右没有明显延迟增长。最主要的好处是服务端线程不再被无意义的轮询请求占满,同样的机器规模可以撑住更多用户。
要做到这个效果有两个前提:一是 Nginx 和网关的缓冲必须关掉,二是事件发送线程池不能有阻塞操作。只要这两点做到了,SSE 在中小规模集群上性能完全够用。
6.3 从SSE中枢到事件总线的演进方向
这套方案目前的定位是“轻量事件中枢”,适合 Agent 并发量在几百到千级的阶段。再往上走,有几个可以升级的方向:
- 把 Redis Pub/Sub 替换成 Redis Stream 的消费者组,获得消费确认、消费者重平衡和持久化能力,事件不丢。
- 事件量特别大时,再用 Kafka 或 Pulsar 这类专业消息队列,把思考事件、工具事件、审计日志分 topic 管理。
- 基于这套中枢扩展 Agent 记忆:把每次任务的工具结果、思考关键片段写入 Redis,带 TTL 做短期记忆,或者定期归并到向量库做长期记忆。
我在实际落地中的体会是:SSE + Redis 这套组合的定位,不是去替代大而全的消息中间件,而是用最小的成本把 Agent 的流式体验做出来。等技术演进到需要更强的事件保证时,接口层面已经隔离好了,内部从 Pub/Sub 切到 Stream 或者其他中间件,下游客户端感知不到变化。对大多数做 Agent 应用的团队来说,先用这套方案把产品体验跑通,比一开始就上一套重中间件要划算得多。