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

Spring WebFlux集成Google Vertex AI流式API实现SSE时,响应块无法实时推送问题咨询

Spring WebFlux集成Google Vertex AI流式API实现SSE时,响应块无法实时推送问题咨询

看起来你遇到的核心问题是阻塞式的流遍历破坏了Spring WebFlux的非阻塞特性,导致所有SSE chunk要等整个流完成才会一次性推送给客户端。我来帮你拆解问题并给出修复方案:

问题根源分析

你在Flux.create的回调里用了同步阻塞的for循环遍历Vertex AI的ResponseStream:

for (GenerateContentResponse response : stream) {
    // 处理chunk...
    sink.next(chunk);
}

Vertex AI的ResponseStream的遍历是阻塞式的——当调用stream()或者直接for循环时,SDK会阻塞当前线程直到下一个响应chunk到来,直到所有响应完成才会退出循环。

而Spring WebFlux的Flux.create回调默认运行在WebFlux的IO事件循环线程上,阻塞这个线程会导致:

  1. 整个反应式流被卡住,直到所有chunk都被遍历完成
  2. 所有sink.next(chunk)的调用都会被缓存,直到循环结束才会一次性推送给客户端

这就是为什么你能在日志里看到每个chunk被处理,但客户端要等全部完成才收到数据的原因。

修复方案:用反应式操作符替代阻塞遍历

我们需要把Vertex AI的阻塞式ResponseStream转换成非阻塞的反应式流,同时把阻塞操作隔离到专门的线程池,避免影响WebFlux的事件循环。

1. 修改Service层的流处理逻辑

把原来的Flux.create替换成Reactor原生的操作符,利用subscribeOn将阻塞遍历放到弹性线程池:

public Flux<String> analyzeFileWithPromptStream(Path filePath, String mimeType, DocumentType documentType, String prompt) {
    try {
        // ... 前面的contentRequest构建逻辑保持不变 ...

        // 替代原来的Flux.create部分
        return Flux.fromIterable(() -> generativeModel.generateContentStream(contentRequest).stream())
                // 将阻塞的流遍历操作放到弹性线程池,避免阻塞WebFlux IO线程
                .subscribeOn(Schedulers.boundedElastic())
                // 展开每个响应的Candidate列表
                .flatMap(response -> Flux.fromIterable(response.getCandidatesList()))
                // 展开每个Candidate的Part列表
                .flatMap(candidate -> Flux.fromIterable(candidate.getContent().getPartsList()))
                // 提取文本内容
                .map(Part::getText)
                // 清理格式
                .map(text -> text.replaceAll("```json", "").replaceAll("```", "").trim())
                // 过滤空chunk
                .filter(chunk -> !chunk.isEmpty())
                .doOnNext(chunk -> log.info("Streaming chunk: {}", chunk));
    } catch (IOException e) {
        log.error("Error reading file: {}", e.getMessage(), e);
        return Flux.error(new RuntimeException("Error reading file", e));
    }
}

2. 确保控制器返回正确的SSE响应头

在你的控制器方法上明确指定响应的媒体类型为text/event-stream,避免客户端误解响应格式:

@PostMapping(value = "/v3/stream", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public Flux<ServerSentEvent<String>> performOcrStream(
        @RequestPart("file") MultipartFile file,
        @RequestPart("request") OcrReqDto request,
        @RequestPart(value = "prompt", required = false) String prompt) {
    // ... 现有逻辑保持不变 ...
}

关键细节解释

  • Flux.fromIterable(() -> ...stream()):用Supplier包装stream()调用,延迟流的初始化,避免提前阻塞线程。
  • subscribeOn(Schedulers.boundedElastic()):这是核心!boundedElastic是Reactor专门为阻塞操作设计的线程池,会自动管理线程数量,把阻塞的流遍历操作和WebFlux的非阻塞IO线程隔离开,确保事件循环不被阻塞。
  • 链式操作替代嵌套循环:用flatMap展开嵌套的Candidate和Part列表,让每个chunk都成为反应式流中的独立元素,这样一旦有新的chunk到来,就能立即推送给客户端。

验证方式

你之前测试的Flux.interval能正常实时推送,是因为它是纯非阻塞的反应式流。修复后,你应该能在Postman或浏览器里看到chunk实时逐个出现,而不是等全部完成才显示。

内容来源于stack exchange

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.07 06:58:05