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事件循环线程上,阻塞这个线程会导致:
- 整个反应式流被卡住,直到所有chunk都被遍历完成
- 所有
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
相关产品推荐
相关产品推荐

