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

如何轻松将JSON流转换为NDJSON流?

JSON 响应流转 NDJSON 流实现方案

核心思路是通过自定义 TransformStream 处理 fetch 返回的字节流,将 JSON 格式(数组或单个对象)实时拆分为符合 NDJSON 规范的单行数据,再传递给 ndjsonStream 处理。

自定义转换流实现

这个流会处理分块传输的 JSON 数据,拼接缓冲区内容并拆分出单个 JSON 元素,输出为带换行的 NDJSON 格式:

class JsonToNdjsonTransform extends TransformStream {
  constructor() {
    super({
      start(controller) {
        this.buffer = '';
        this.isArray = false;
        this.depth = 0;
      },
      transform(chunk, controller) {
        this.buffer += chunk;
        
        let index = 0;
        while (index < this.buffer.length) {
          const char = this.buffer[index];

          if (!this.isArray) {
            // 检测是否为数组开头
            if (char === '[') {
              this.isArray = true;
              index++;
              continue;
            }
            // 处理单个 JSON 对象
            if (char === '{' || char === '[') this.depth++;
            if (char === '}' || char === ']') this.depth--;
            
            if (this.depth === 0 && (char === '}' || char === ']')) {
              const element = this.buffer.slice(0, index + 1).trim();
              controller.enqueue(element + '\n');
              this.buffer = this.buffer.slice(index + 1);
              index = 0;
              continue;
            }
            index++;
          } else {
            // 处理数组内的元素
            if (char === '{' || char === '[') this.depth++;
            if (char === '}' || char === ']') this.depth--;

            if (char === '}' && this.depth === 1) {
              const elementEnd = index + 1;
              const element = this.buffer.slice(0, elementEnd).trim();
              if (element) controller.enqueue(element + '\n');
              
              this.buffer = this.buffer.slice(elementEnd);
              index = 0;
              // 跳过元素间的逗号和空白
              while (this.buffer[0] === ',' || /\s/.test(this.buffer[0])) {
                this.buffer = this.buffer.slice(1);
              }
            } else if (this.depth === 0) {
              // 数组结束,清空缓冲区
              this.buffer = '';
              index = this.buffer.length;
            } else {
              index++;
            }
          }
        }
      },
      flush(controller) {
        // 处理剩余的合法 JSON 内容
        const remaining = this.buffer.trim();
        if (remaining) {
          try {
            JSON.parse(remaining);
            controller.enqueue(remaining + '\n');
          } catch (e) {
            console.error('剩余内容不是合法 JSON:', e);
          }
        }
      }
    });
  }
}

使用方式

结合 fetch 和 pipeThrough 串联流,将 JSON 流转为 NDJSON 流后再交给 ndjsonStream:

const response = await fetch(url);

// 处理压缩的响应(如果需要)
const decompressedStream = response.body.pipeThrough(new DecompressionStream('gzip'));

// 转换为 NDJSON 流
const ndjsonStream = decompressedStream
  .pipeThrough(new TextDecoderStream())
  .pipeThrough(new JsonToNdjsonTransform())
  .pipeThrough(new TextEncoderStream());

// 交给 ndjsonStream 处理
const ndjsonStreamer = ndjsonStream(ndjsonStream).getReader();

关键细节

  • 分块处理:通过缓冲区拼接不完整的 JSON 片段,确保实时处理,无需等待整个响应下载完成
  • 格式兼容:同时支持 JSON 数组和单个 JSON 对象的转换
  • 错误容忍:flush 阶段验证剩余内容的合法性,避免输出无效数据
  • 压缩支持:如果响应是 gzip 压缩,需先通过 DecompressionStream 解压

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.23 13:25:26