ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Spring Boot 3 + SSE + Redis:生产级 AI Agent 流式输出与工具调用实战

Spring Boot 3 + SSE + Redis:生产级 AI Agent 流式输出与工具调用实战 1. 为什么流式 Agent 的“最后一公里”总是翻车做过 AI Agent 项目的人大概率都遇到过这种场景前端页面已经渲染出了思考过程的骨架用户盯着屏幕等第一个 token 落地结果等了十几秒浏览器控制台突然蹦出一行stream disconnected before completion: idle timeout waiting for sse然后整个对话窗口卡死用户只能刷新重来。更尴尬的是后端日志显示模型其实早就把结果吐完了只是中间某个环节把数据“吞”了。这个问题的根源往往不在大模型本身而在流式传输链路的设计。Agent 的流式输出和普通聊天机器人的流式输出有本质区别普通对话是“一问一答”的单次流而 Agent 是“思考—工具调用—再思考—再调用—最终回答”的多阶段流。中间可能穿插数据库查询、外部 API 调用、文件读写每个阶段都可能产生几十秒的空窗期。如果传输层用的是轮询前端就得不停发请求问“好了没”如果用的是短连接 SSE空闲超时就会把连接掐断。我这次要拆解的这套方案核心就是用Spring Boot 3 SSE Redis搭一个生产级的 Agent 流式思考与工具调用中枢。它要解决三件事第一让思考过程和工具调用结果都能实时推到前端不轮询第二用 Redis 做跨节点的状态中转和分布式协调保证多实例部署下流不丢、不重、不乱序第三把空闲超时、断流重连、工具调用阻塞这些坑提前填掉。适合谁看如果你正在做 AI Agent 产品、需要把 Agent 的执行过程可视化、或者被 SSE 断流和 Redis 超时折磨过这篇内容基本可以照着抄。我会把选型理由、参数计算、代码骨架、排查表都摊开讲尽量做到“看完就能落地”。2. 整体架构设计与技术选型拆解2.1 为什么是 SSE 而不是 WebSocket流式传输可选的技术路线主要有三条轮询、WebSocket、SSE。轮询的问题最明显Agent 一次执行可能持续 30 秒到几分钟轮询间隔设短了浪费请求设长了用户体验差而且每次轮询都要重新建立上下文服务端压力成倍增加。WebSocket 是全双工理论上更适合交互但它有两个现实问题。一是运维复杂度高Nginx、网关、负载均衡对 WebSocket 的长连接支持需要额外配置连接数一多内存和文件描述符消耗很可观。二是 Agent 流式输出本质上是服务端单向推送用户在前端只需要“接收”不需要在同一个连接上频繁回传数据WebSocket 的双工能力用不上反而增加了心跳维护和断线重连的复杂度。SSE 基于 HTTP 协议天然支持单向推送浏览器端用EventSource就能接服务端在 Spring Boot 3 里用SseEmitter或者 WebFlux 的FluxServerSentEvent都能实现。它的优势是走标准 HTTP网关和负载均衡友好自动重连机制内置文本协议调试时直接 curl 就能看到流。对于 Agent 这种“服务端持续推、客户端持续收”的场景SSE 是性价比最高的选择。注意SSE 默认是文本协议如果工具调用结果里有二进制内容需要先 Base64 编码再推送否则会破坏事件流格式。2.2 Redis 在链路里到底扮演什么角色很多人第一反应是SSE 不是长连接吗为什么还要 Redis单实例部署确实可以不用但生产环境不可能只跑一个节点。一旦多实例部署就会出现三个问题。第一连接粘性问题。用户 A 的 SSE 连接落在节点 1但 Agent 执行任务可能被调度到节点 2节点 2 产生的思考流怎么推给节点 1 上的连接答案是通过 Redis 的 Pub/Sub 做跨节点广播节点 2 把流事件发布到 Redis 频道所有节点订阅持有该用户连接的节点负责推送。第二状态共享问题。Agent 的执行状态当前处于哪个阶段、已调用哪些工具、中间结果是什么需要跨节点可见否则断线重连后无法恢复上下文。这里用 Redis 的 Hash 结构存执行快照用 String 结构存幂等键用 List 做事件缓冲。第三分布式锁问题。同一个会话如果被重复触发可能导致 Agent 重复执行、工具重复调用。用 Redis 分布式锁保证同一会话同一时刻只有一个执行实例锁的过期时间要大于 Agent 最大执行时间同时配合看门狗续期。热词里出现的redis command timed out; nested exception is io.lettuce.core.RedisCommandTimeoutException就是典型的 Redis 超时问题后面排查章节会专门讲。2.3 整体数据流设计整套链路的数据流可以拆成四段前端发起请求建立 SSE 连接服务端返回SseEmitter同时生成一个streamId。请求进入 Agent 执行器执行器把streamId和会话信息写入 Redis并获取分布式锁。Agent 每产生一个事件思考片段、工具调用开始、工具调用结果、最终回答就封装成事件对象发布到 Redis 频道agent:stream:{streamId}。持有 SSE 连接的节点订阅该频道收到事件后通过SseEmitter.send()推给前端同时把事件追加到 Redis List 做缓冲支持断线重连后的补发。这个设计的关键在于执行与推送解耦。Agent 执行器不关心谁在监听只管往 Redis 发SSE 节点不关心谁在执行只管从 Redis 收。两边通过 Redis 解耦多实例部署时天然支持水平扩展。3. 核心细节解析与实操要点3.1 SSE 连接的生命周期管理SseEmitter有三个关键回调onCompletion、onTimeout、onError。很多人只设置了超时时间却忘了在回调里清理资源导致连接泄漏。超时时间的设置有个计算逻辑。Agent 单次执行的最大时长假设是 5 分钟那么 SSE 超时不能设成 5 分钟因为工具调用阶段可能有空窗期。我的经验值是SSE 超时 Agent 最大执行时长 × 1.5 30 秒缓冲。按 5 分钟算就是 480 秒。同时前端EventSource本身也有重连机制服务端超时后前端会自动重连重连时带上streamId服务端从 Redis List 里补发未收到的事件。SseEmitter emitter new SseEmitter(480_000L); emitter.onCompletion(() - { log.info(SSE completed, streamId{}, streamId); redisTemplate.opsForHash().delete(agent:stream:meta, streamId); }); emitter.onTimeout(() - { log.warn(SSE timeout, streamId{}, streamId); emitter.complete(); }); emitter.onError((e) - { log.error(SSE error, streamId{}, streamId, e); emitter.completeWithError(e); });实操心得onCompletion里不要做耗时操作否则会阻塞容器线程。清理 Redis 的动作建议丢到异步线程池里执行。3.2 Redis 频道订阅与事件分发Redis Pub/Sub 的订阅要用独立的连接不能和业务操作的连接池混用。Spring Data Redis 里通过RedisMessageListenerContainer配置每个订阅者持有独立连接。事件对象的设计要包含几个字段streamId、eventTypethinking / tool_start / tool_result / answer / done、sequence序号用于排序和去重、payload内容、timestamp。sequence很关键因为 Redis Pub/Sub 不保证跨频道顺序多节点并发发布时可能乱序前端按sequence排序后再渲染。public record AgentStreamEvent( String streamId, String eventType, long sequence, String payload, long timestamp ) {}分发逻辑是订阅者收到事件后先判断本地是否有对应的SseEmitter有就直接推没有就说明连接不在本节点忽略即可其他节点会处理。同时无论本地有没有连接都要把事件追加到 Redis Listagent:stream:buffer:{streamId}并设置过期时间防止内存泄漏。3.3 分布式锁与幂等控制分布式锁用 Redis 的SET key value NX PX timeout实现value 用 UUID 保证只有加锁者能解锁。锁的 key 设计为agent:lock:{sessionId}过期时间设为 Agent 最大执行时长 60 秒。String lockKey agent:lock: sessionId; String lockValue UUID.randomUUID().toString(); Boolean locked redisTemplate.opsForValue() .setIfAbsent(lockKey, lockValue, Duration.ofSeconds(360)); if (Boolean.FALSE.equals(locked)) { throw new BizException(该会话正在执行中请稍后); }解锁时用 Lua 脚本保证原子性先比对 value 再删除避免误删别人的锁。if redis.call(get, KEYS[1]) ARGV[1] then return redis.call(del, KEYS[1]) else return 0 end注意锁的过期时间一定要大于 Agent 最大执行时间。如果 Agent 执行超过锁过期时间锁自动释放另一个请求就能进来导致重复执行。稳妥做法是加一个看门狗线程每 30 秒续期一次。3.4 工具调用阶段的空窗期处理Agent 调用工具时可能几十秒没有输出这时候 SSE 连接虽然没断但前端会以为卡死了。解决办法是心跳事件。服务端每隔 15 秒往流里推一个eventTypeheartbeat的空事件前端收到后刷新“正在执行”的状态提示同时保持连接活跃。心跳的实现可以放在 Agent 执行器的独立线程里也可以利用 Redis 的过期监听。我倾向于用ScheduledExecutorService定时往 Redis 频道发心跳简单可控。scheduler.scheduleAtFixedRate(() - { if (isExecuting(streamId)) { publishEvent(streamId, heartbeat, ); } }, 0, 15, TimeUnit.SECONDS);4. 实操过程与核心环节实现4.1 环境准备与依赖配置先确认版本Spring Boot 3.2.x、Spring Data Redis 3.2.x、Lettuce 6.3.x。Lettuce 是 Spring Boot 3 默认的 Redis 客户端基于 Netty支持异步和响应式比 Jedis 更适合高并发场景。pom.xml核心依赖dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-web/artifactId /dependency dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-data-redis/artifactId /dependency dependency groupIdorg.apache.commons/groupId artifactIdcommons-pool2/artifactId /dependencyapplication.yml里 Redis 配置要重点关注超时和连接池spring: data: redis: host: 127.0.0.1 port: 6379 timeout: 3000ms lettuce: pool: max-active: 16 max-idle: 8 min-idle: 2 max-wait: 2000mstimeout: 3000ms是命令超时不是连接超时。热词里那个RedisCommandTimeoutException很多时候就是这个值设太小Agent 执行期间 Redis 操作被阻塞导致超时。建议设 3000ms 以上同时排查是否有大 key 或慢查询。4.2 SSE 接口的完整实现Controller 层暴露一个GET /agent/stream接口接收sessionId和query参数返回SseEmitter。GetMapping(value /agent/stream, produces MediaType.TEXT_EVENT_STREAM_VALUE) public SseEmitter stream(RequestParam String sessionId, RequestParam String query) { String streamId UUID.randomUUID().toString(); SseEmitter emitter new SseEmitter(480_000L); emitterRegistry.register(streamId, emitter); agentExecutor.executeAsync(streamId, sessionId, query); return emitter; }emitterRegistry是一个本地ConcurrentHashMapString, SseEmitter存本节点持有的连接。Agent 执行器是异步的通过Async或者线程池提交不阻塞 HTTP 线程。4.3 Agent 执行器的流式事件发布执行器内部按阶段推进每个阶段产生事件后立即发布到 Redis。public void executeAsync(String streamId, String sessionId, String query) { String lockKey agent:lock: sessionId; String lockValue UUID.randomUUID().toString(); if (!tryLock(lockKey, lockValue)) { publishEvent(streamId, error, 会话正在执行中); return; } try { long seq 0; publishEvent(streamId, thinking, 开始分析问题, seq); // 阶段一意图识别 String intent llmClient.classify(query); publishEvent(streamId, thinking, 识别意图 intent, seq); // 阶段二工具调用 ListToolCall tools llmClient.planTools(intent); for (ToolCall tool : tools) { publishEvent(streamId, tool_start, tool.getName(), seq); String result toolExecutor.execute(tool); publishEvent(streamId, tool_result, result, seq); } // 阶段三生成最终回答 String answer llmClient.generate(query, tools); publishEvent(streamId, answer, answer, seq); publishEvent(streamId, done, , seq); } finally { unlock(lockKey, lockValue); } }publishEvent内部做两件事往 Redis 频道发消息往 Redis List 追加缓冲。private void publishEvent(String streamId, String type, String payload, long seq) { AgentStreamEvent event new AgentStreamEvent( streamId, type, seq, payload, System.currentTimeMillis()); String json objectMapper.writeValueAsString(event); redisTemplate.convertAndSend(agent:stream: streamId, json); redisTemplate.opsForList().rightPush(agent:stream:buffer: streamId, json); redisTemplate.expire(agent:stream:buffer: streamId, Duration.ofMinutes(10)); }4.4 订阅端的事件推送订阅端配置RedisMessageListenerContainer监听所有agent:stream:*频道。Bean public RedisMessageListenerContainer listenerContainer( RedisConnectionFactory factory, AgentStreamSubscriber subscriber) { RedisMessageListenerContainer container new RedisMessageListenerContainer(); container.setConnectionFactory(factory); container.addMessageListener(subscriber, new PatternTopic(agent:stream:*)); return container; }Subscriber 收到消息后解析出streamId从本地注册表找SseEmitter找到就推。Override public void onMessage(Message message, byte[] pattern) { String json new String(message.getBody()); AgentStreamEvent event objectMapper.readValue(json, AgentStreamEvent.class); SseEmitter emitter emitterRegistry.get(event.streamId()); if (emitter ! null) { try { emitter.send(SseEmitter.event() .id(String.valueOf(event.sequence())) .name(event.eventType()) .data(event.payload())); } catch (IOException e) { emitterRegistry.remove(event.streamId()); } } }4.5 断线重连与事件补发前端EventSource断线后会自动重连重连时带上Last-Event-ID头值是最后收到的事件sequence。服务端在建立新连接时从 Redis List 里读取sequence大于该值的事件依次补发。GetMapping(value /agent/stream/reconnect, produces MediaType.TEXT_EVENT_STREAM_VALUE) public SseEmitter reconnect(RequestParam String streamId, RequestHeader(value Last-Event-ID, required false) String lastId) { SseEmitter emitter new SseEmitter(480_000L); emitterRegistry.register(streamId, emitter); if (lastId ! null) { long lastSeq Long.parseLong(lastId); ListString buffer redisTemplate.opsForList() .range(agent:stream:buffer: streamId, 0, -1); for (String json : buffer) { AgentStreamEvent event objectMapper.readValue(json, AgentStreamEvent.class); if (event.sequence() lastSeq) { emitter.send(SseEmitter.event() .id(String.valueOf(event.sequence())) .name(event.eventType()) .data(event.payload())); } } } return emitter; }实操心得补发时要注意顺序Redis List 是插入顺序但多节点并发发布时可能乱序补发前按sequence排序一次更稳妥。5. 常见问题与排查技巧实录5.1 空闲超时断流的三种排查方向stream disconnected before completion: idle timeout waiting for sse这个报错排查要分三层。第一层是网关层。Nginx 默认proxy_read_timeout是 60 秒如果 Agent 执行超过 60 秒没有数据推送Nginx 会主动断开。解决办法是在 Nginx 配置里把proxy_read_timeout和proxy_send_timeout调到 600 秒以上同时关闭proxy_buffering否则 Nginx 会缓冲 SSE 数据导致前端收不到实时流。location /agent/stream { proxy_pass http://backend; proxy_read_timeout 600s; proxy_send_timeout 600s; proxy_buffering off; proxy_cache off; proxy_set_header Connection ; proxy_http_version 1.1; chunked_transfer_encoding off; }第二层是服务端 SSE 超时。SseEmitter构造时传入的超时时间如果小于 Agent 执行时间到点就会触发onTimeout。按前面算的 480 秒设置同时配合心跳事件基本不会误触发。第三层是前端 EventSource 超时。浏览器对 SSE 连接也有默认超时虽然大多数现代浏览器不会主动断但如果服务端长时间不发数据某些代理或浏览器插件会掐断。心跳事件就是解决这个的。5.2 Redis 命令超时的定位与解决redis command timed out; nested exception is io.lettuce.core.RedisCommandTimeoutException这个报错常见原因有四个。一是大 key 操作。比如agent:stream:buffer:{streamId}这个 List 如果一直追加不清理可能涨到几万条range操作就会很慢。解决办法是设置过期时间同时限制 List 长度用LTRIM只保留最近 1000 条。二是慢查询阻塞。Redis 是单线程处理命令如果有一个KEYS *或者大SMEMBERS在执行后续命令都会排队。用SLOWLOG GET 10查看慢查询把slowlog-log-slower-than设成 10000 微秒。三是连接池耗尽。max-active设太小高并发时拿不到连接就会超时。按 QPS 估算假设峰值 500 QPS每个操作平均 2ms理论上 1 个连接就够但实际要考虑网络抖动建议max-active设为 16 到 32。四是网络抖动。跨机房访问 Redis 时延迟高timeout设 3000ms 可能不够。这种情况要么把 Redis 和应用部署在同机房要么调大超时时间。5.3 分布式锁失效导致重复执行分布式锁失效的典型表现是同一个会话被触发两次Agent 执行了两遍工具被调用了两次用户看到重复的回答。排查思路先看锁的过期时间是否小于 Agent 最大执行时间。如果 Agent 执行了 6 分钟锁只设了 5 分钟第 5 分钟锁自动释放第二个请求就能进来。解决办法是加看门狗续期或者把锁过期时间设得足够大。另一个原因是解锁时误删。如果 A 线程加锁后执行超时锁自动释放B 线程加锁成功此时 A 线程执行完去解锁如果直接DEL就会把 B 的锁删掉。所以解锁必须用 Lua 脚本比对 value。5.4 事件乱序与重复推送Redis Pub/Sub 不保证顺序多节点并发发布时前端可能先收到tool_result再收到tool_start。解决办法是前端按sequence排序后再渲染服务端保证sequence单调递增。重复推送的原因是断线重连时补发逻辑和实时推送重叠。比如补发到第 50 条时实时流又推了第 50 条。解决办法是前端按sequence去重维护一个已处理的最大sequence小于等于该值的事件直接丢弃。5.5 常见问题速查表问题现象可能原因排查方法解决方案SSE 空闲超时断流网关超时、无心跳查 Nginx 日志、看心跳间隔调大 proxy_read_timeout、加心跳Redis 命令超时大 key、慢查询、连接池小SLOWLOG、监控连接数LTRIM 限长、调大连接池锁失效重复执行锁过期时间短、误删锁看执行时长、查解锁日志看门狗续期、Lua 解锁事件乱序Pub/Sub 无序看 sequence 字段前端按 sequence 排序事件重复补发与实时重叠看 sequence 重复前端按 sequence 去重连接泄漏未清理注册表看内存和连接数onCompletion 里异步清理避坑技巧生产环境一定要给agent:stream:buffer:{streamId}设过期时间否则 Redis 内存会被慢慢吃光。我一般设 10 分钟足够覆盖断线重连窗口。6. 性能压测与容量估算6.1 单节点能扛多少 SSE 连接SSE 连接是长连接每个连接占用一个 HTTP 连接和一个SseEmitter对象。Tomcat 默认最大连接数是 8192但实际能扛多少取决于内存和文件描述符。按经验值每个 SSE 连接约占 10KB 到 50KB 内存取决于缓冲区大小一台 4GB 内存的机器预留 2GB 给 JVM 堆大概能扛 4 万个连接。但文件描述符是更硬的限制Linux 默认ulimit -n是 1024需要调到 65535。ulimit -n 65535Tomcat 的maxConnections也要相应调大server: tomcat: max-connections: 20000 accept-count: 1000 threads: max: 500 min-spare: 50注意threads.max不需要设太大因为 SSE 连接建立后就释放了工作线程只有推送时才短暂占用。6.2 Redis Pub/Sub 的吞吐瓶颈Redis Pub/Sub 的吞吐很高单实例能到 10 万 QPS 以上。但 Agent 场景下每个事件都要发布一次如果 Agent 每秒产生 10 个事件1000 个并发 Agent 就是 1 万 QPS单实例 Redis 完全扛得住。真正的瓶颈在订阅端。每个节点都要订阅所有频道收到消息后要判断本地有没有连接。如果节点数多每个节点都要处理全量消息浪费 CPU。优化办法是用 Redis 的Sharded Pub/SubRedis 7.0按streamId分片每个节点只订阅一部分频道。6.3 容量估算示例假设业务规模日活 1 万峰值并发 Agent 执行 500 个每个 Agent 平均执行 60 秒每秒产生 5 个事件。SSE 连接数500 个单节点扛得住。事件发布 QPS500 × 5 2500 QPSRedis 单实例轻松应对。事件缓冲内存每个事件约 500 字节每个 Agent 产生 300 个事件500 个 Agent 就是 75MB设 10 分钟过期峰值内存约 150MB。锁数量500 个每个锁约 100 字节可忽略。按这个规模一台 4 核 8GB 的应用节点 一台 2 核 4GB 的 Redis 节点就能撑住。如果要水平扩展应用节点无状态加节点即可Redis 如果成为瓶颈可以升级到集群模式。7. 一些踩过的坑和收尾经验这套方案我在两个项目里落地过踩过的坑不少挑几个最有代表性的说。第一个坑是心跳事件把前端搞崩了。一开始心跳事件也走answer类型前端把它当正文渲染结果页面上出现一堆空行。后来给心跳单独定义了heartbeat类型前端收到后只更新状态提示不渲染内容。第二个坑是Redis List 没限长。有个 Agent 执行了 20 分钟产生了上万条事件agent:stream:buffer涨到 5MB补发时range操作直接把 Redis 阻塞了。后来加了LTRIM只保留最近 1000 条同时把过期时间从 30 分钟调到 10 分钟。第三个坑是锁的续期线程没关。看门狗线程在 Agent 执行完后没有正确关闭导致线程泄漏跑了一天之后线程数涨到几千。后来用ScheduledFuture持有任务句柄执行完在finally里cancel。第四个坑是Nginx 缓冲导致流不实时。默认proxy_buffering onNginx 会攒一批数据再发前端看到的就是一段一段的不是逐字输出。关掉proxy_buffering后恢复正常。最后分享一个小技巧调试 SSE 的时候直接用 curl 最直观。curl -N -H Accept: text/event-stream http://localhost:8080/agent/stream?sessionIdtestqueryhello-N参数关闭 curl 的缓冲能实时看到服务端推的每一条事件。比开浏览器 DevTools 方便得多尤其是排查事件格式和顺序问题时。这套东西后续还能扩展的方向不少比如把事件缓冲从 Redis List 换成 Redis Stream支持多消费者组和更精确的 ACK 机制或者把 Agent 执行器做成独立服务通过消息队列和 SSE 节点通信进一步解耦。但就目前这套 Spring Boot 3 SSE Redis 的组合已经能覆盖大多数生产场景了。
返回列表