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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 12:18:17