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
相关产品推荐
相关产品推荐

