如何正确解析按换行分隔的文本流中的JSON对象
问题描述
我通过Fetch API读取返回流中的JSON对象,希望将解析出的对象存入数组后传入转换器执行自定义转换,但当前代码无法处理分片产生的不完整JSON块,导致部分对象解析失败,寻求完善方案或更优实现方式。
原异步迭代器代码
function streamAsyncIterator(stream) { const reader = stream.pipeThrough(new TextDecoderStream()).getReader(); return { next() { return reader.read(); }, return() { reader.releaseLock(); return {}; }, [Symbol.asyncIterator]() { return this; } }; }
原解析逻辑代码
const fetchProductsStream = async () => { const body = { searchRange: ['fragrance', 'perfumery', 'cologne'], exclusionCategory: ['discounted'] } let time1 = performance.now(); const response = await fetch('/products/stream', { method: 'POST', body: JSON.stringify(body) }); let chunkCount = 0; let remaining = ""; let currentChunk = ""; let totalItems = 0; for await (const chunk of streamAsyncIterator(response.body)) { let items = []; try{ currentChunk = chunk; let lastIndex; if(remaining.trim().length > 1){ currentChunk = remaining + chunk; } currentChunk.split(/\r?\n/).forEach((line, index) => { try{ items.push(JSON.parse(line)); console.log(`Pushed line ${line} to items array`); lastIndex = index; }catch(err){ console.error(`Error is ${err}`); } }) totalItems += items.length; remaining = ""; let i = currentChunk.indexOf(lastIndex); if(i != -1){ remaining = currentChunk.substring(lastIndex); } items = []; }catch(error){ console.log(error); } chunkCount++; } let time2 = performance.now(); console.log(`Total number of items is ${totalItems}`); console.log(`Total number of chunks :: ${chunkCount} in time ${(time2 - time1)/1000} seconds`) }
示例数据
{"productId":"1671","category":"fragrance","categoryId":"1WW23","productName":"Paco Rabane","productClass":"M","warehouseCode":"1242","alternateName":"ALT_ELL322"} {"productId":"1671","category":"fragrance","categoryId":"1WW23","productName":"Paco Rabane","productClass":"M","warehouseCode":"1242","alternateName":"ALT_ELL322"} {"productId":"1671","category":"fragrance","categoryId":"1WW23","productName":"Paco Rabane","productClass":"M","warehouseCode":"1242","alternateName":"ALT_ELL322","category":{"desc":"1654 Cologne Perfume","code":"221","feature":"1","salesCode":"S2237"}} {"productId":"1671","category":"fragrance","categoryId":"1WW23","productName":"Paco Rabane","productClass":"M","warehouseCode":"1242","alternateName":"ALT_ELL322"}
修复方案与更优实现
问题分析
原代码核心问题在于剩余不完整JSON块的处理逻辑错误:
- 错误使用数组索引截取字符串,而非保留分割后最后一个不完整行
- 未正确区分可解析的完整行与需留存的不完整行
修复后的解析逻辑
const fetchProductsStream = async () => { const body = { searchRange: ['fragrance', 'perfumery', 'cologne'], exclusionCategory: ['discounted'] } const time1 = performance.now(); const response = await fetch('/products/stream', { method: 'POST', body: JSON.stringify(body) }); let chunkCount = 0; let remaining = ""; let totalItems = 0; for await (const chunk of streamAsyncIterator(response.body)) { chunkCount++; // 拼接剩余内容与当前分片 const fullText = remaining + chunk; // 按换行分割内容 const lines = fullText.split(/\r?\n/); // 最后一行可能不完整,留存到下一轮处理 remaining = lines.pop() || ""; // 处理所有完整行 for (const line of lines) { if (!line.trim()) continue; // 跳过空行 try { const item = JSON.parse(line); totalItems++; // 此处可直接传入转换器执行自定义转换 // transformItem(item); console.log(`解析成功:${JSON.stringify(item)}`); } catch (err) { console.error(`解析行失败:${err.message},内容:${line}`); } } } // 处理流结束后剩余的最后一段内容 if (remaining.trim()) { try { const item = JSON.parse(remaining); totalItems++; // transformItem(item); console.log(`解析剩余内容成功:${JSON.stringify(item)}`); } catch (err) { console.error(`解析剩余内容失败:${err.message},内容:${remaining}`); } } const time2 = performance.now(); console.log(`总解析条目数:${totalItems}`); console.log(`总块数:${chunkCount},耗时:${(time2 - time1)/1000}秒`) }
更优实现:使用TransformStream处理流
利用浏览器原生TransformStream可更优雅地处理流转换,无需手动管理迭代器与剩余内容:
async function fetchProductsStream() { const body = { searchRange: ['fragrance', 'perfumery', 'cologne'], exclusionCategory: ['discounted'] }; const time1 = performance.now(); const response = await fetch('/products/stream', { method: 'POST', body: JSON.stringify(body) }); let remaining = ""; let totalItems = 0; let chunkCount = 0; // 创建转换流:将文本分片解析为JSON对象 const parserStream = new TransformStream({ transform(chunk, controller) { chunkCount++; const text = new TextDecoder().decode(chunk); const fullText = remaining + text; const lines = fullText.split(/\r?\n/); remaining = lines.pop() || ""; for (const line of lines) { if (!line.trim()) continue; try { const item = JSON.parse(line); totalItems++; // 将解析后的对象发送到下游,可直接串联自定义转换流 controller.enqueue(item); } catch (err) { console.error(`解析失败:${err.message},内容:${line}`); } } }, flush(controller) { // 处理流结束后剩余的最后一段内容 if (remaining.trim()) { try { const item = JSON.parse(remaining); totalItems++; controller.enqueue(item); } catch (err) { console.error(`解析剩余内容失败:${err.message},内容:${remaining}`); } } } }); // 处理转换后的流数据 const reader = response.body.pipeThrough(parserStream).getReader(); while (true) { const { done, value } = await reader.read(); if (done) break; // 此处执行自定义转换逻辑 // const transformed = transformItem(value); console.log(`获取到解析后的对象:${JSON.stringify(value)}`); } const time2 = performance.now(); console.log(`总解析条目数:${totalItems}`); console.log(`总块数:${chunkCount},耗时:${(time2 - time1)/1000}秒`); }
优势
- 贴合Stream API设计规范,代码结构更模块化
- 无需手动实现异步迭代器,借助原生API简化逻辑
- 可直接串联自定义转换流,实现解析+转换的流水线处理
内容的提问来源于stack exchange,提问作者BreenDeen
相关产品推荐
相关产品推荐

