如何确保StreamedResponse的Chunk为完整字符串以避免JSON解码异常
问题描述
我的服务器会发送JSON格式的Server-side Event(SSE)数据。在“加载”阶段,负载数据较短,但当流完成时,我会遇到FormatException,原因是负载超过了Chunk的content-length。
我希望在将Chunk传入jsonDecode前确保它是完整的字符串,是否可行?
初始客户端Dart代码
final imageEventProvider = StreamProvider.autoDispose<Map<String, dynamic>>( (ref) async* { final request = http.Request( 'POST', backendUrl().replace( pathSegments: ['gen-image'], ), ); final auth = ref.read(authProvider); request.body = jsonEncode({'prompt': 'a desert landscape'}); request.headers["Accept"] = "text/event-stream"; request.headers["Cache-Control"] = "no-cache"; request.headers['Content-Type'] = 'application/json'; request.headers['X-FIREBASE-TOKEN'] = await auth.currentUser!.getIdToken(); final sResp = await request.send(); final predictionStream = sResp.stream.transform(utf8.decoder); await for (final part in predictionStream) { final Map<String, dynamic> data = jsonDecode(part); if (data['status'] != 'succeeded') { ref.state = const AsyncLoading(); } yield data; } }, );
补充信息
服务器端(Python)
事件模型代码:
@dataclass class SSEvent: message: dict[str, Any] def encode(self) -> bytes: return json.dumps(self.message).encode("utf-8")
响应结构代码:
async def sse(): pred_resp = await client.post( url, json={"version": version, "input": input}, headers=headers, ) data = pred_resp.json() pred_url = "{}/{}".format(url, data["id"]) while data["status"] != "succeeded": pred_resp = await client.post(pred_url, headers=headers) data = pred_resp.json() event = SSEvent({"status": "loading"}) yield event.encode() else: data = pred_resp.json() event = SSEvent( { **data, } ) yield event.encode() response = await make_response( sse(), { "Content-Type": "text/event-stream", "Cache-Control": "no-cache", "Transfer-Encoding": "chunked", }, )
修改后的客户端(Dart)
按照建议修改后仍出现FormatException:
final imageEventProvider = StreamProvider.autoDispose<Map<String, dynamic>>( (ref) async* { final request = http.Request( 'POST', backendUrl().replace( pathSegments: ['gen-image'], ), ); final auth = ref.read(authProvider); request.body = jsonEncode({'prompt': 'a desert landscape'}); request.headers["Accept"] = "text/event-stream"; request.headers["Cache-Control"] = "no-cache"; request.headers['Content-Type'] = 'application/json'; request.headers['X-FIREBASE-TOKEN'] = await auth.currentUser!.getIdToken(); final sResp = await request.send(); final predictionStream = sResp.stream.transform(json.fuse(utf8).decoder); final data = await predictionStream.single as Map<String, dynamic>; if (data!['status'] != 'succeeded') { ref.state = const AsyncLoading(); } yield data; }, );
解决方案
问题核心在于服务器未遵循标准SSE格式,且客户端未正确处理SSE流的事件分割逻辑。
1. 修复服务器端SSE格式
标准SSE要求每个事件必须以data: 开头,以两个换行符\n\n作为事件分隔符,客户端才能正确识别独立事件。当前服务器直接返回JSON字节,无SSE格式标记,导致流中多个JSON事件混在一起,大JSON还会被拆分为多个Chunk,直接解析必然报错。
修改Python的SSEvent.encode方法:
def encode(self) -> bytes: json_str = json.dumps(self.message) # 按SSE标准包装事件内容 return f"data: {json_str}\n\n".encode("utf-8")
2. 客户端正确处理SSE流
客户端需要先按SSE格式分割事件,再提取每个事件中的JSON内容,而非直接解码整个流。可以通过StreamTransformer实现事件分割逻辑:
final imageEventProvider = StreamProvider.autoDispose<Map<String, dynamic>>( (ref) async* { final request = http.Request( 'POST', backendUrl().replace( pathSegments: ['gen-image'], ), ); final auth = ref.read(authProvider); request.body = jsonEncode({'prompt': 'a desert landscape'}); request.headers["Accept"] = "text/event-stream"; request.headers["Cache-Control"] = "no-cache"; request.headers['Content-Type'] = 'application/json'; request.headers['X-FIREBASE-TOKEN'] = await auth.currentUser!.getIdToken(); final sResp = await request.send(); // 处理SSE流:分割事件、提取JSON内容 final sseStream = sResp.stream .transform(utf8.decoder) .transform(StreamTransformer<String, String>.fromHandlers( handleData: (data, sink) { // 按SSE事件分隔符\n\n拆分 final events = data.split('\n\n'); for (final event in events) { if (event.isNotEmpty) { // 提取data: 前缀后的JSON字符串 final jsonStr = event.replaceFirst('data: ', ''); sink.add(jsonStr); } } }, )); await for (final jsonStr in sseStream) { if (jsonStr.isEmpty) continue; try { final Map<String, dynamic> data = jsonDecode(jsonStr); if (data['status'] != 'succeeded') { ref.state = const AsyncLoading(); } yield data; } catch (e) { // 捕获解析错误,比如不完整的Chunk片段 print('JSON解析失败: $e'); } } }, );
为什么之前的修改无效?
你之前使用的json.fuse(utf8).decoder是将整个字节流当作单个JSON对象解析,但服务器返回的是多个SSE事件(多个独立JSON),且大JSON可能被拆分为多个Chunk,这种方式完全不适用SSE场景。
额外说明
- 若服务器Chunk分割导致单个JSON被拆分为多段,上述客户端逻辑会自动缓存不完整片段,直到收到完整的
\n\n分隔符才会解析,确保处理完整的JSON字符串。 - 生产环境可考虑使用成熟的SSE客户端库(如
sse_client),简化流分割和事件处理的细节。
内容的提问来源于stack exchange,提问作者Mudassir
相关产品推荐
相关产品推荐

