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

Java Spring中如何等待TEXT_EVENT_STREAM订阅完成并同步返回结果?

解决WebClient调用SSE流同步拼接字符串的问题

问题分析

你原代码用subscribe()异步收集数据,方法返回时数据还未写入List,所以得到空结果;直接用blockLast()报错,是因为Reactor的IO线程(如reactor-http-nio-3)禁止阻塞操作——这类线程专为非阻塞IO设计,阻塞会导致线程池耗尽,破坏异步调度策略。

解决方案

情况1:方法运行在允许阻塞的线程(如Spring MVC控制器线程)

直接通过Reactor操作符收集并拼接数据,最后调用block()获取同步结果:

private String executeStream(String streamId, UserToken currentUserAndToken) {
    URI streamUri = UriComponentsBuilder.newInstance()
            .scheme("https")
            .host(host)
            .path(streamPath)
            .buildAndExpand(streamVariablesMap)
            .toUri();

    return webClient.get()
            .uri(streamUri)
            .headers(h -> h.setBearerAuth(currentUserAndToken.getToken()))
            .accept(MediaType.TEXT_EVENT_STREAM)
            .retrieve()
            .bodyToFlux(String.class)
            .collectList() // 收集所有流元素到List
            .map(list -> String.join(", ", list)) // 拼接成目标字符串
            .block(); // 阻塞直到流结束并返回结果
}

情况2:方法运行在Reactor非阻塞线程(如WebFlux控制器、过滤器线程)

需要先切换到支持阻塞操作的线程池(Schedulers.boundedElastic()),再调用block():

private String executeStream(String streamId, UserToken currentUserAndToken) {
    URI streamUri = UriComponentsBuilder.newInstance()
            .scheme("https")
            .host(host)
            .path(streamPath)
            .buildAndExpand(streamVariablesMap)
            .toUri();

    return webClient.get()
            .uri(streamUri)
            .headers(h -> h.setBearerAuth(currentUserAndToken.getToken()))
            .accept(MediaType.TEXT_EVENT_STREAM)
            .retrieve()
            .bodyToFlux(String.class)
            .collectList()
            .map(list -> String.join(", ", list))
            .publishOn(Schedulers.boundedElastic()) // 切换到支持阻塞的线程池
            .block();
}

关键说明

  • collectList()会等待整个流结束后收集所有元素,仅适用于有限流;如果是无限流,需要用take(n)限制元素数量,或者改用reduce实时拼接:reduce("", (acc, data) -> acc + (acc.isEmpty() ? "" : ", ") + data)。
  • Schedulers.boundedElastic()是Reactor专为阻塞操作提供的线程池,会自动控制线程数量,避免资源耗尽。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 15:30:13