Spring Boot MVC RestClient是否支持Flux类型SSE流?Servlet项目无WebFlux的Flux客户端
Spring Boot MVC中SSE流处理的两个技术问题解答
背景
在Spring WebFlux中,我们可以借助Flux类型和反应式WebClient优雅处理SSE数据,示例代码如下:
@GetMapping(value = "/chat", produces = MediaType.TEXT_EVENT_STREAM_VALUE) public Flux<String> streamChatResponses() { // 示例:调用OpenAI接口 return webClient.post() .uri("/chat/completions") .contentType(MediaType.APPLICATION_JSON) .bodyValue(requestBody) .retrieve() .bodyToFlux(String.class); }
对于Spring Boot MVC(Servlet)项目,我们可以引入Spring AI让控制器返回Flux类型的SSE数据,但一直缺少支持Flux类型的HTTP客户端:
<dependency> <groupId>com.alibaba.cloud.ai</groupId> <artifactId>spring-ai-alibaba-starter</artifactId> </dependency>
@RestController public class SseController { @GetMapping(value = "/sse", produces = MediaType.TEXT_EVENT_STREAM_VALUE) public Flux<String> streamEvents() { // TODO: Spring Boot MVC的RestClient或RestTemplate是否支持返回Flux类型的SSE? return Flux.interval(Duration.ofSeconds(1)) .map(sequence -> "SSE Event - " + LocalTime.now()); } }
原始的SSE流处理实现较为基础,需要手动处理JSON转换和响应设置:
@GetMapping("/sse-stream-old-code") public void streamSse(HttpServletResponse response) { response.setContentType("text/event-stream"); response.setCharacterEncoding(StandardCharsets.UTF_8.name()); RestTemplate restTemplate = new RestTemplate(); String url = "http://*******/chat_stream"; ResponseExtractor<Void> responseExtractor = restResponse -> { try (BufferedReader reader = new BufferedReader(new InputStreamReader(restResponse.getBody(), StandardCharsets.UTF_8))) { String line; while ((line = reader.readLine()) != null) { response.getWriter().write("data: " + line + "\n\n"); response.getWriter().flush(); } } return null; }; restTemplate.execute(url, HttpMethod.POST, null, responseExtractor); }
技术问题解答
1. Spring Boot MVC的RestClient是否支持返回Flux类型的SSE流数据?
RestClient是Spring 6+推出的同步HTTP客户端,本身不原生支持反应式流(如Flux),但可以通过以下方式实现将SSE流转换为Flux:
- 通过RestClient的
exchange方法获取响应的InputStream; - 使用Reactor的
Flux.create工具类,将输入流的读取逻辑封装为Flux,实现异步流式处理。
示例代码如下:
import reactor.core.publisher.Flux; @GetMapping(value = "/sse-restclient", produces = MediaType.TEXT_EVENT_STREAM_VALUE) public Flux<String> streamWithRestClient() { return Flux.create(sink -> { RestClient restClient = RestClient.create(); restClient.post() .uri("http://*******/chat_stream") .contentType(MediaType.APPLICATION_JSON) .bodyValue(requestBody) .exchange((request, response) -> { try (BufferedReader reader = new BufferedReader( new InputStreamReader(response.getBody(), StandardCharsets.UTF_8))) { String line; while ((line = reader.readLine()) != null) { if (sink.isCancelled()) { break; } sink.next("data: " + line + "\n\n"); } sink.complete(); } catch (IOException e) { sink.error(e); } return null; }); }); }
2. 在Spring Boot(Servlet)项目中,是否存在无需引入WebFlux依赖即可支持Flux数据流的客户端?
存在。只需引入reactor-core依赖(无需完整的WebFlux),结合普通的HTTP客户端(如OkHttp、Apache HttpClient),即可手动封装出支持Flux的流式HTTP客户端。
方案1:使用OkHttp + reactor-core
首先添加依赖:
<dependency> <groupId>io.projectreactor</groupId> <artifactId>reactor-core</artifactId> </dependency> <dependency> <groupId>com.squareup.okhttp3</groupId> <artifactId>okhttp</artifactId> </dependency>
示例实现:
import okhttp3.OkHttpClient; import okhttp3.Request; import okhttp3.RequestBody; import okhttp3.Response; import okhttp3.ResponseBody; import reactor.core.publisher.Flux; @GetMapping(value = "/sse-okhttp", produces = MediaType.TEXT_EVENT_STREAM_VALUE) public Flux<String> streamWithOkHttp() { OkHttpClient client = new OkHttpClient(); Request request = new Request.Builder() .url("http://*******/chat_stream") .post(RequestBody.create(requestBody.getBytes(), okhttp3.MediaType.parse("application/json"))) .build(); return Flux.create(sink -> { client.newCall(request).enqueue(new okhttp3.Callback() { @Override public void onFailure(okhttp3.Call call, IOException e) { sink.error(e); } @Override public void onResponse(okhttp3.Call call, Response response) throws IOException { try (ResponseBody body = response.body()) { if (body == null) { sink.complete(); return; } BufferedReader reader = new BufferedReader(body.charStream()); String line; while ((line = reader.readLine()) != null) { if (sink.isCancelled()) { break; } sink.next("data: " + line + "\n\n"); } sink.complete(); } catch (IOException e) { sink.error(e); } } }); }); }
方案2:使用RestClient + reactor-core
无需额外引入HTTP客户端,仅依赖Spring的RestClient和reactor-core,实现方式同问题1中的示例,无需引入WebFlux依赖。
内容的提问来源于stack exchange,提问作者zy_sun
相关产品推荐
相关产品推荐

