如何轻松将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
相关产品推荐
相关产品推荐

