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

使用Observable无法分块获取全部数据的问题排查

为什么RxJS fromFetch没实现流式数据传输?

你当前代码的问题出在response.json()这个方法上——它会等待整个响应体完全加载后,一次性解析成完整的JSON对象,所以Observable只会触发一次next回调,把整个数组传过来,根本做不到分块流式传输。哪怕服务器是分块发送数据,response.json()也会等所有块接收完毕才处理。

要实现真正的流式传输,你得直接操作响应的ReadableStream,逐块读取并解析数据。如果服务器返回的是JSON Lines格式(每行一个独立的JSON对象),可以用下面的代码修改:

const fetch = (url, payload) => {
  try {
    const requestOptions = {
      method: 'POST',
      headers: { 'Content-Type': 'application/json' },
      body: JSON.stringify(payload),
    };
    const request = new Request(url, requestOptions);
    const data$ = fromFetch(request).pipe(
      switchMap(response => {
        if (!response.ok) {
          return of({ error: true, message: `Error ${response.status}` });
        }
        // 直接操作ReadableStream实现流式读取
        const reader = response.body.getReader();
        const decoder = new TextDecoder();
        let buffer = '';

        return new Observable(observer => {
          const processChunk = async () => {
            const { done, value } = await reader.read();
            if (done) {
              // 处理最后剩余的不完整行
              if (buffer.trim()) {
                try {
                  observer.next(JSON.parse(buffer));
                } catch (e) {
                  observer.error(e);
                }
              }
              observer.complete();
              return;
            }
            // 解码二进制数据并拼接缓冲区
            buffer += decoder.decode(value, { stream: true });
            // 按换行分割成独立行
            const lines = buffer.split('\n');
            buffer = lines.pop(); // 把未完成的行留在缓冲区
            // 逐行解析JSON并发送
            for (const line of lines) {
              if (line.trim()) {
                try {
                  observer.next(JSON.parse(line));
                } catch (e) {
                  observer.error(e);
                }
              }
            }
            processChunk();
          };
          processChunk().catch(err => observer.error(err));
          // 取消订阅时终止流读取
          return () => reader.cancel();
        });
      }),
      catchError(err => {
        console.error(err);
        return of({ error: true, message: err.message });
      })
    );

    data$.subscribe({
      next: result => console.log('收到单条数据:', result),
      error: err => console.error('流式错误:', err),
      complete: () => console.log('流式传输完成')
    });
  } catch (e) {
    console.error(e);
  }
};

补充说明

  • 服务器需要配合返回JSON Lines格式:每个数据对象单独占一行,而不是包裹在一个大数组里。如果服务器还是返回分块的完整数组,你得自己处理不完整JSON片段的解析(复杂度会高很多)。
  • 通过response.body.getReader()获取流读取器,逐块解码后按行分割,就能把每个独立的JSON对象逐个通过observer.next()发送,真正实现分块流式传输。

内容的提问来源于stack exchange,提问作者BreenDeen

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.09 02:35:22