SpringBoot代理Flask SSE端点出现分片不均问题求助
Flask SSE端点经SpringBoot代理后分片异常问题解决
我有一个符合MDN标准的Flask SSE端点,返回的SSE数据分片以两个\n字符结尾,前端直接调用时能收到完整分片。但通过SpringBoot的SseEmitter实现代理后,前端调用代理端点时出现分片错乱,比如部分分片仅收到data:,实际事件值在后续分片里。以下是相关代码及测试结果,求解决方法:
示例Flask应用
import json import time from flask import Flask, Response from flask_cors import CORS, cross_origin app = Flask(__name__) cors = CORS( app, origins="*", ) def markdown_event(): yield "data: "+ json.dumps({"data": "{\"RESPONSE\": \"well_structured_markdown_formatted_response_with_citation_without_appending_references_and_sources_and_additional_questions\", \"HEADING\": \"\", \"NEXT_HEADING\": \"suggestive_questions\", \"REFERENCE\": \"id_list\", \"IS_SORRY\": \"\", \"IS_RELEVANT\": \"\", \"LANGUAGE_OF_QUERY\": \"\", \"TEXT_CLASS\": \"\"}", "type": "schema"}) + "\n\n" chunks = ["{\"", "well", "_struct", "ured", "_mark", "down", "_formatted", "_response", "_with", "_c", "itation", "_without", "_app", "ending", "_references", "_and", "_sources", "_and", "_additional", "_questions", "\":", " \"", "Hello", " there", " [", "2", "]", " how", " can", " i", " be", " sure", " that", " [", "1", ",", "2", "]", " citation", " are", " being", " covered", " [", "1", " -", " ", "4", "],", " more", " citation", " formats", ":", " [", "1", ",", " ", "3", " -", " ", "5", ".", " Some", " more", " [", "1", "][", "4", "](", "3", ".", " The", " years", " are", " like", " [", "199", "5", " -", " ", "199", "8", "]", " []", " String", ":", " '", "Hello", " there", " [", "4", "]", " how", " can", " i", " be", " sure", " that", " (", "1", ",", "2", ")", " citation", " are", " being", ".", " Hello", " [", "202", "1", "]", " the", " ", " college", " year", " was", " [", "199", "2", "-", "96", ".", "\\", "n", "\\n", "The", " task", " is", " simple", " you", " should", " remove", " citations", " from", " the", " random", " above", " string", " in", " format", " like", " [", "xxx", "]", " [", "xxx", " -", " xxx", "],", " [", "xxx", " ,", " xxx", "],", " where", " every", " x", " is", " a", " digit", ";", " and", " not", " years", " in", " the", " following", " formats", " [", "xxxx", "],", " [", "xxxx", " -", " xxxx", "],", " [", "xxxx", " -", " xx", "\\", "n", "\\n", "The", " output", " of", " the", " above", " string", " would", " be", ":\\", "n", "Hello", " there", " ", " how", " can", " i", " be", " sure", " that", " ", " citation", " are", " being", " covered", " ,", " more", " citation", " formats", ":", " .", " Some", " more", " .", " The", " years", " are", " like", " [", "199", "5", " -", " ", "199", "8", "]", " []", " String", ":", " '", "Hello", " there", " ", " how", " can", " i", " be", " sure", " that", " ", " citation", " are", " being", ".", " Hello", " [", "202", "1", "]", " the", " ", " college", " year", " was", " [", "199", "2", "-", "96", ".", "\\", "n", "\\n", "See", " that", " was", " easy", ".\"", " \"", "reference", "_ids", "_mapping", "\":"," [{\"", "citation", "_number", "\":", " ", "1", ",", " \"", "reference", "_id", "\":", " \"", "01", "HT", "G", "AA", "28", "DB", "H", "8", "A", "1", "ND", "P", "XY", "RK", "SW", "GE", "\"},", " {\"", "citation", "_number", "\":", " ", "2", ",", " \"", "reference", "_id", "\":", " \"", "01", "HP", "7", "CV", "V", "38", "J", "BC", "7", "D", "RA", "4", "EW", "QM", "FS", "71", "\"},", " {\"", "citation", "_number", "\":", " ", "3", ",", " \"", "reference", "_id", "\":", " \"", "01", "HP", "72", "M", "DS", "19", "J", "TF", "J", "Y", "WF", "6", "G", "5", "A", "AY", "7", "W", "\"}", "]", " ", "suggest", "ive", "_questions", "\":", " [\"", "What", " are", " the", " benefits", " of", " Q", "IV", " in", " children", " aged", " ", "6", " months", " to", " ", "17", " years", "?\",", " \"", "Can", " Q", "IV", " be", " administered", " to", " children", " aged", " ", "6", " months", " and", " above", "?\",", " \"", "What", " is", " the", " safety", " profile", " of", " Q", "IV", " in", " children", " aged", " ", "6", "-", "35", " months", "?", "\"}"] for chunk in chunks: yield "data: "+json.dumps({"data": chunk, "type":"function_call"}) + "\n\n" time.sleep(0.1) @app.route("/") def index(): return "Hello, go to /markdown endpoint" @app.route("/markdown", methods=["GET", "POST"]) def markdown(): return Response(markdown_event(), mimetype="text/event-stream") if __name__ == "__main__": app.run(port=5500)
启动Flask应用命令
python -m flask --app=app run --port=5500 --host=0.0.0.0
示例前端客户端
async function sse() { var requestOptions = { method: "POST", redirect: "follow", }; const response = await fetch( "http://localhost:5000/proxy-sse", // "http://localhost:5500/markdown", requestOptions ); console.log(response.status) console.log(response.headers) count = 1 const reader = response.body .pipeThrough(new TextDecoderStream()) .getReader(); while (true) { const { value, done } = await reader.read(); if (done) break; console.log("Chunk: ", count, " *********************************"); const a = String.raw`${value}`; console.log("Received", a); console.log("--------------------------------"); count += 1; } console.log(count); } sse()
SpringBoot控制器代码(原异常代码)
package com.spring.demo.controller; import java.io.IOException; import java.time.LocalDateTime; import org.springframework.http.MediaType; import org.springframework.web.bind.annotation.PostMapping; import org.springframework.web.bind.annotation.RestController; import org.springframework.web.reactive.function.client.WebClient; import org.springframework.web.servlet.mvc.method.annotation.SseEmitter; import org.springframework.web.servlet.mvc.method.annotation.SseEmitter.SseEventBuilder; @RestController public class HelloWorldController { private final String targetSseUrl = "http://localhost:5500/"; private final WebClient webClient; public HelloWorldController(WebClient.Builder webClientBuilder) { this.webClient = webClientBuilder.baseUrl(targetSseUrl).build(); } @PostMapping("/proxy-sse") public SseEmitter proxySse() { SseEmitter emitter = new SseEmitter(Long.MAX_VALUE); webClient.post() .uri("/markdown") .accept(MediaType.TEXT_EVENT_STREAM) .retrieve() .bodyToFlux(String.class) .subscribe(arg0 -> { try { System.err.println(arg0); emitter.send(SseEmitter.event().data(arg0)); } catch (IOException e) { e.printStackTrace(); } }, emitter::completeWithError, emitter::complete); return emitter; } @PostMapping("/local-sse") public SseEmitter proxySse1() { SseEmitter emitter = new SseEmitter(Long.MAX_VALUE); emitter.onTimeout(() -> System.out.println("SSE connection timed out")); emitter.onCompletion(() -> System.out.println("SSE connection closed")); new Thread(() -> { int count = 0; while (true) { count++; try { String data = "Event at: " + LocalDateTime.now(); SseEventBuilder eventBuilder = SseEmitter.event().data(data); emitter.send(eventBuilder); Thread.sleep(100); } catch (Exception e) { emitter.completeWithError(e); break; } if (count == 100) { emitter.complete(); break; } } }).start(); return emitter; } }
各端点测试结果
- Flask端点:前端收到的每个分片都是完整的SSE事件格式(
data: xxx\n\n),无拆分情况,能正常解析。 - SpringBoot代理SSE端点:出现分片错乱,部分分片仅包含
data:,对应的数据内容在后续分片返回,导致前端无法正确识别单个SSE事件。 - SpringBoot本地SSE端点:本地生成的SSE事件能完整传递,前端收到的分片正常,无异常。
问题根源
原代码中使用WebClient的bodyToFlux(String.class)时,WebClient会根据底层缓冲区大小自动拆分数据流,而不是按照SSE的事件分隔符\n\n进行分割,导致单个完整的SSE事件被拆分成多个片段发送给前端。
修复方案
自定义数据流处理逻辑,按SSE的事件分隔符\n\n分割完整事件后再转发,修改后的SpringBoot代理代码如下:
package com.spring.demo.controller; import java.io.IOException; import java.nio.charset.StandardCharsets; import org.springframework.core.io.buffer.DataBuffer; import org.springframework.core.io.buffer.DataBufferUtils; import org.springframework.http.MediaType; import org.springframework.web.bind.annotation.PostMapping; import org.springframework.web.bind.annotation.RestController; import org.springframework.web.reactive.function.client.WebClient; import org.springframework.web.servlet.mvc.method.annotation.SseEmitter; import reactor.core.publisher.Flux; @RestController public class HelloWorldController { private final String targetSseUrl = "http://localhost:5500/"; private final WebClient webClient; public HelloWorldController(WebClient.Builder webClientBuilder) { this.webClient = webClientBuilder.baseUrl(targetSseUrl).build(); } @PostMapping("/proxy-sse") public SseEmitter proxySse() { SseEmitter emitter = new SseEmitter(Long.MAX_VALUE); webClient.post() .uri("/markdown") .accept(MediaType.TEXT_EVENT_STREAM) .retrieve() .bodyToFlux(DataBuffer.class) .transform(this::splitBySseDelimiter) .subscribe( event -> { try { emitter.send(event); } catch (IOException e) { emitter.completeWithError(e); } }, emitter::completeWithError, emitter::complete ); return emitter; } private Flux<String> splitBySseDelimiter(Flux<DataBuffer> dataBuffers) { return DataBufferUtils.join(dataBuffers) .map(buffer -> { String content = StandardCharsets.UTF_8.decode(buffer.asByteBuffer()).toString(); DataBufferUtils.release(buffer); return content; }) .flatMapMany(content -> Flux.fromArray(content.split("\n\n"))) .filter(event -> !event.trim().isEmpty()) .map(event -> event + "\n\n"); } @PostMapping("/local-sse") public SseEmitter proxySse1() { SseEmitter emitter = new SseEmitter(Long.MAX_VALUE); emitter.onTimeout(() -> System.out.println("SSE connection timed out")); emitter.onCompletion(() -> System.out.println("SSE connection closed")); new Thread(() -> { int count = 0; while (true) { count++; try { String data = "Event at: " + java.time.LocalDateTime.now(); emitter.send(SseEmitter.event().data(data)); Thread.sleep(100); } catch (Exception e) { emitter.completeWithError(e); break; } if (count == 100) { emitter.complete(); break; } } }).start(); return emitter; } }
修复原理
- 读取原始
DataBuffer数据流,避免WebClient自动拆分字符串。 - 将所有
DataBuffer拼接成完整字符串,再按SSE标准的事件分隔符\n\n分割成单个完整事件。 - 过滤分割后产生的空字符串,确保每个事件都是有效内容。
- 为每个事件重新添加
\n\n,保证转发的数据流符合SSE格式,前端能正确识别每个事件。
内容的提问来源于stack exchange,提问作者Cdaman
相关产品推荐
相关产品推荐

