You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用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;
}

额外排查点

  1. 确认上游API的SSE流是否正常结束:可以用curl测试上游接口,看是否在发送[DONE]后会关闭连接,如果上游流一直不关闭,Flux的doOnComplete永远不会触发,自然无法调用sseEmitter.complete()
  2. 检查是否有未捕获的异常:在订阅者或doOnError中添加详细日志,确认是否有异常导致complete()未被执行

内容的提问来源于stack exchange,提问作者Cloud Ja

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.08 06:16:08