Node.js中如何使用stream-json的pipeline实现文件写入操作
实现思路梳理
- stream-chain 提供的
chain()作用是将多个流、转换函数拼接为一个独立的管道流,你代码中将文件读取流、解压流作为chain的前两个参数是完全正确的,chain本身已经包含了数据源,不需要额外引入dataSource。 stream-json的parser()输出的是JSON语法标记流(比如对象起始标记、键名、值片段等结构化数据),而非可直接写入文件的JSON文本,如果不需要修改JSON内容,完全可以跳过parser环节,直接将解压后的流写入文件。如果需要处理JSON内容,需要在parser后搭配对应的处理组件,再将处理结果序列化为文本才能写入目标文件。
基础实现方案
场景1:仅解压写入文件,无需处理JSON内容
不需要引入stream-json相关组件,直接拼接流即可:
const { chain } = require('stream-chain'); const fs = require('fs'); const zlib = require('zlib'); const Path = require('path'); function unzipJson() { const zipPath = Path.resolve(__dirname, 'resources', 'myfile.json.zip'); const jsonPath = Path.resolve(__dirname, 'resources', 'myfile.json'); console.info('Attempting to read zip'); return new Promise((resolve, reject) => { const pipeline = chain([ fs.createReadStream(zipPath), zlib.createGunzip(), fs.createWriteStream(jsonPath) ]); // 监听流核心事件即可 pipeline.on('end', () => { console.info('Unzip and write completed'); resolve(); }); pipeline.on('error', (err) => { console.error('Process failed: ', err); reject(err); }); }); }
场景2:需要处理JSON内容后再写入
需要搭配stream-json的处理组件,再通过自定义转换流将处理后的JSON数据序列化为文本:
const { chain } = require('stream-chain'); const { parser } = require('stream-json'); const { streamValues } = require('stream-json/streamers/StreamValues'); const fs = require('fs'); const zlib = require('zlib'); const Path = require('path'); const { Transform } = require('stream'); // 自定义转换流:将处理后的JSON对象转为可写入的文本格式 const jsonStringifyTransform = new Transform({ writableObjectMode: true, transform(chunk, _, callback) { // chunk为streamValues输出的{key: xxx, value: xxx}结构 this.push(JSON.stringify(chunk.value) + '\n'); callback(); } }); function unzipAndProcessJson() { const zipPath = Path.resolve(__dirname, 'resources', 'myfile.json.zip'); const jsonPath = Path.resolve(__dirname, 'resources', 'processed.json'); return new Promise((resolve, reject) => { const pipeline = chain([ fs.createReadStream(zipPath), zlib.createGunzip(), parser(), // 此处可插入你需要的过滤组件:比如pick、ignore,或者自定义转换函数 streamValues(), // 示例自定义过滤逻辑:仅保留符合条件的数据 data => { const value = data.value; return value?.department === 'accounting' ? data : null; }, jsonStringifyTransform, fs.createWriteStream(jsonPath) ]); pipeline.on('end', resolve); pipeline.on('error', reject); }); }
事件监听说明
你需要关注的核心事件只有两个:
end事件:流处理全流程完成、所有数据都写入目标文件后触发,用来执行完成后的回调逻辑error事件:处理流程中任意环节报错都会触发,用来捕获处理异常
data事件是每次有数据片段流过时触发,如果你不需要逐段统计、打印进度这类需求,不需要监听,写入流会自动处理接收到的所有数据。
内容的提问来源于stack exchange,提问作者AncientSwordRage
相关产品推荐
相关产品推荐

