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

如何确保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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.10 10:54:54