ARTICLE DETAIL

资讯详情

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

Spring AI 流式调用过程中的异常捕获与断网重连机制

Spring AI 流式调用过程中的异常捕获与断网重连机制 Spring AI 流式调用过程中的异常捕获与断网重连机制与传统微服务之间毫秒级的短 HTTP 交互不同基于大模型LLM的流式生成Streaming Response调用链路通常会持续10 秒到数分钟之久。在如此漫长的长连接生命周期中遇到公网网络抖动、机房交换机短暂丢包、上游供应商接口超时SocketTimeoutException或连接意外被重置PrematureCloseException几乎是不可避免的常态。在非流式场景下处理异常非常简单配置一个重试拦截器重新发起请求即可。然而在流式场景下如果一个长回答已经在前端打字输出了 800 字在第 801 字时连接发生异常断开直接整单重试会导致前端页面内容被清空重刷不仅造成用户体验崩塌还会白白浪费前 800 字的 Token 费用与推理时间。为了保障流式交互的工业级健壮性我们需要在 Spring AI 与响应式流Reactor Flux体系中构建一套区分阶段的异常捕获、流中熔断兜底与断点增量补全机制。流式异常的双阶段划分与处理策略流式调用的异常必须严格划分为两大阶段采取完全不同的恢复策略┌─────────────────────────┐ │ 发起 LLM 流式调用请求 │ └───────────┬─────────────┘ │ 遇到网络/服务异常? │ ┌────────────────────────┴────────────────────────┐ ▼ (尚未接收到任何 Token) ▼ (流传输过程中截断) 【阶段一连接前置异常】 【阶段二流中截断异常】 - 上游 429 / 503 / 鉴权失败 - ReadTimeout / PrematureClose - 触发全局退避重试 (Retry.backoff) - 严禁从头重试 - 或透明切换备用模型渠道 - 发起增量 Continuation 补偿请求阶段一前置握手与首字前异常此时上游尚未吐出任何有效 Token可安全利用 Reactor 的retryWhen结合指数退避Exponential Backoff重新发起请求或直接 Failover 切换至备用机房/模型账号。阶段二流中传输截断Mid-Stream Failure此时必须保留已经下发到缓冲区的历史文本基于已生成的截断内容构造“请基于前文继续补充”的增量续写 Prompt实现无缝断点续接。核心实现基于 Reactor 的流式容错管道利用 Project Reactor 提供的onErrorResume、doOnError与原子状态机在 Spring AI 的流处理流中实现精细化异常拦截与兜底package com.example.ai.stream; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.ai.chat.client.ChatClient; import org.springframework.stereotype.Service; import reactor.core.publisher.Flux; import reactor.util.retry.Retry; import java.time.Duration; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicReference; Service public class ResilientStreamChatService { private static final Logger log LoggerFactory.getLogger(ResilientStreamChatService.class); private final ChatClient chatClient; public ResilientStreamChatService(ChatClient.Builder builder) { this.chatClient builder.build(); } public FluxString streamWithResilience(String userPrompt) { StringBuilder receivedContent new StringBuilder(); AtomicBoolean hasStartedReceiving new AtomicBoolean(false); AtomicReferenceString lastValidChunk new AtomicReference(); return chatClient.prompt() .user(userPrompt) .stream() .content() .doOnNext(chunk - { hasStartedReceiving.set(true); receivedContent.append(chunk); lastValidChunk.set(chunk); }) // 1. 前置连接异常重试仅在尚未收到任何 Chunk 时允许整单重试 .retryWhen(Retry.backoff(2, Duration.ofMillis(500)) .filter(throwable - !hasStartedReceiving.get()) .doBeforeRetry(retrySignal - log.warn(前置连接异常正在发起第 {} 次重试..., retrySignal.totalRetries() 1)) ) // 2. 流中异常捕获与降级补偿 .onErrorResume(throwable - { log.error(流式推理过程中断! 当前已接收 {} 字符, 异常类型: {}, receivedContent.length(), throwable.getClass().getSimpleName(), throwable); if (!hasStartedReceiving.get()) { // 连首字都没出来就彻底挂了返回友好错误提示 return Flux.just(\n[系统繁忙未能连接至模型服务请稍后重试]); } // 已经输出了部分内容尝试发起增量 Continuation 补偿续写 return attemptContinuation(userPrompt, receivedContent.toString()); }); } private FluxString attemptContinuation(String originalPrompt, String partialGeneratedText) { log.info(正在尝试增量断点续接...); String continuationPrompt String.format( 你刚才在回答以下问题\n\%s\\n\n在生成到以下内容时连接意外中断\n\%s\\n\n请直接紧接着上述末尾内容继续输出不要重复前面已经生成过的任何字词。, originalPrompt, // 截取末尾 200 字提供精准断点上下文 partialGeneratedText.substring(Math.max(0, partialGeneratedText.length() - 200)) ); return chatClient.prompt() .user(continuationPrompt) .stream() .content() .onErrorReturn(\n[网络连接波动输出已截断]); } }协议层断线重连结合 SSE Last-Event-ID 机制除了后端的主动补偿前端与网关之间应充分利用 Server-Sent EventsSSE的原生重连协议[前端 EventSource] ──────── 接收到 Event-ID: 42 ────────► [网络中断] │ ▼ (浏览器原生自动重连) [发起重连 HTTP 请求] ─── 带上请求头: Last-Event-ID: 42 ───► [API Gateway] │ ▼ [从 Redis 环形缓冲区回放 43 之后的帧] ◄───────────────────────┘服务端事件帧缓存设计环形缓冲区Ring Buffer在后端网关或业务节点中为每个活跃会话分配一个轻量级的环形内存缓存保留最近 60 秒内下发的最后 100 个 SSE 事件帧。断点回放当网关捕获到带有Last-Event-ID: 42的重连请求时首先从环形缓冲区中提取 ID 42 的存量数据进行快速重放然后再将客户端重新绑定到活跃的响应式流管道中。生产落地的安全防线防止无限续接递归在增量补偿逻辑attemptContinuation中必须严格将重试深度限制为1 次。如果补偿请求再次发生网络中断应立即终止并输出终止标记坚决避免陷入死循环重试导致计费失控。客户端超时保护前端在建立流式连接时应设置两级超时首字超时First Token Timeout建议设为 15 秒若 15 秒内未收到任何 SSE 帧主动断开并提示重试。流中静默超时Idle Timeout若在流传输过程中超过 8 秒未收到任何新 Token 或心跳 Ping判定为假死连接主动触发 Abort 并启动重连流程。收益总结在某企业级智能知识库助手的生产调优中通过引入“前置退避重试 流中增量续接”将公网弱网环境下用户的对话中断失败率从8.4% 压降至 0.3%大幅提升了流式产品的稳定品质。
返回列表