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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.04 07:54:04