使用Flux.subscribe结合SseEmitter时连接无法关闭的问题
SseEmitter调用complete()后连接无法关闭的问题排查与解决
你遇到的问题核心是WebClient返回的Flux订阅线程与SseEmitter的响应处理线程上下文不一致,或者异步发送操作未完成就调用了complete(),导致Spring MVC无法正确感知SSE流结束,连接一直保持直到超时。
解决方案1:将SSE发送操作整合到Reactor流中(推荐)
把SseEmitter.send()的异步操作转换成Reactor的Mono,让Flux等待所有发送完成后再触发complete(),确保流的生命周期对齐:
修改后的代码如下:
@GetMapping("") public SseEmitter chat(@Validated ChatRequest request) { // ... 其他初始化代码 WebClient webClient = WebClient.create("http://example.com"); SseEmitter sseEmitter = new SseEmitter(); ModelChatRequest modelChatRequest = modelChatRequest(); StringBuilder builder = new StringBuilder(); // 假设builder是你需要的缓存变量 Flux<ServerSentEvent<String>> eventStream = webClient.post() .uri("/v1/chat/completions") .header("Content-Type", "application/json") .accept(MediaType.TEXT_EVENT_STREAM) .body(BodyInserters.fromValue(modelChatRequest)) .retrieve() .bodyToFlux(new ParameterizedTypeReference<ServerSentEvent<String>>() {}); // 将发送操作整合到Flux处理链中 eventStream.flatMap(event -> { try { String data = event.data(); if (!data.equals("[DONE]")) { ModelChatVO vo = new ObjectMapper().readValue(data, new TypeReference<ModelChatVO>() {}); String content = vo.getChoices().get(0).getDelta().getContent(); if (!StringUtils.isBlank(content)) { builder.append(content); // 将send的异步操作转为Mono,让Flux等待发送完成 return Mono.fromFuture(sseEmitter.send(content)); } } return Mono.empty(); } catch (Exception e) { return Mono.error(e); } }) .doOnComplete(() -> { try { sseEmitter.complete(); } catch (Exception e) { sseEmitter.completeWithError(e); } }) .doOnError(t -> { log.error("处理SSE消息失败", t); sseEmitter.completeWithError(t); }) .subscribe(); // 启动订阅 return sseEmitter; }
这样做的好处是:
- 所有
send()操作都会被Reactor流跟踪,确保数据全部发送完毕后才会触发doOnComplete - 异常会被统一处理,避免因单个消息处理失败导致流卡住
解决方案2:手动处理[DONE]事件并等待发送完成
如果上游会明确发送[DONE]作为结束标记,可以在订阅者中直接处理该事件,同时等待send()操作完成后再结束:
@GetMapping("") public SseEmitter chat(@Validated ChatRequest request) { // ... 其他初始化代码 WebClient webClient = WebClient.create("http://example.com"); SseEmitter sseEmitter = new SseEmitter(); ModelChatRequest modelChatRequest = modelChatRequest(); StringBuilder builder = new StringBuilder(); Flux<ServerSentEvent<String>> eventStream = webClient.post() .uri("/v1/chat/completions") .header("Content-Type", "application/json") .accept(MediaType.TEXT_EVENT_STREAM) .body(BodyInserters.fromValue(modelChatRequest)) .retrieve() .bodyToFlux(new ParameterizedTypeReference<ServerSentEvent<String>>() {}); Disposable subscription = eventStream.subscribe(event -> { try { String data = event.data(); if (data.equals("[DONE]")) { // 收到结束标记,立即关闭SSE连接并取消订阅 sseEmitter.complete(); subscription.dispose(); return; } ModelChatVO vo = new ObjectMapper().readValue(data, new TypeReference<ModelChatVO>() {}); String content = vo.getChoices().get(0).getDelta().getContent(); if (!StringUtils.isBlank(content)) { builder.append(content); // 等待发送完成,避免异步发送未结束就关闭连接 sseEmitter.send(content).get(); } } catch (Exception e) { sseEmitter.completeWithError(e); subscription.dispose(); } }); // 兜底:设置超时时间,避免连接无限等待 sseEmitter.onTimeout(() -> { log.warn("SSE连接超时"); sseEmitter.complete(); subscription.dispose(); }); return sseEmitter; }
额外排查点
- 确认上游API的SSE流是否正常结束:可以用curl测试上游接口,看是否在发送
[DONE]后会关闭连接,如果上游流一直不关闭,Flux的doOnComplete永远不会触发,自然无法调用sseEmitter.complete() - 检查是否有未捕获的异常:在订阅者或
doOnError中添加详细日志,确认是否有异常导致complete()未被执行
内容的提问来源于stack exchange,提问作者Cloud Ja
相关产品推荐
相关产品推荐

