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

如何在Node.js中构建完整ETL管道?解决执行顺序与JSON解析问题

Node.js ETL管道脚本修复方案

问题根源

  1. 异步流程无序:解压操作未等待完成就执行后续步骤,导致读取到不完整的JSON文件,触发SyntaxError: Unexpected end of JSON input。
  2. 回调与Promise兼容问题:回调式的文件读取未包装为Promise,无法被Promise.all正确等待。
  3. 变量意外覆盖:转换步骤中错误地将transform_JSON函数覆盖为返回的字符串。
  4. 无效Promise混入:处理非目标文件时未过滤空值,干扰Promise.all的执行逻辑。

修复后的完整代码

const fs = require('fs');
const { promises: { readdir, readFile, writeFile, mkdir } } = require("fs");
const url = require('url');
const zlib = require('zlib');
const { pipeline } = require('stream/promises'); // 用stream.promises简化流操作

const input_dir = `${__dirname}/input`;
const input_unzipped_dir = `${__dirname}/input-unzipped`;
const output_dir = `${__dirname}/output`;

// 获取目录文件列表
async function get_files(dir) {
  return readdir(dir);
}

// Promise风格的JSON读取函数
async function read_json_file(file_path) {
  const file_data = await readFile(file_path, 'utf-8');
  return JSON.parse(file_data);
}

// 转换JSON数据(纯函数,无外部变量依赖)
function transform_JSON(file_data) {
  const u = url.parse(file_data.u);
  const query_map = new Map(Object.entries(file_data.e));

  return {
    timestamp: file_data.ts,
    url_object: {
      domain: u.host,
      path: u.path,
      query_object: query_map,
      hash: u.hash,
    },
    ec: file_data.e,
  };
}

// Promise风格的文件压缩函数
async function gzip_file(input_path, output_path) {
  const readStream = fs.createReadStream(input_path);
  const gzip = zlib.createGzip();
  const writeStream = fs.createWriteStream(output_path);
  await pipeline(readStream, gzip, writeStream);
}

const orchestrate_etl_pipeline = async () => {
  try {
    // 确保必要目录存在
    await Promise.all([
      !fs.existsSync(input_unzipped_dir) && mkdir(input_unzipped_dir),
      !fs.existsSync(output_dir) && mkdir(output_dir)
    ]);

    // 步骤1:解压.gz文件
    const gz_files = (await get_files(input_dir)).filter(filename => filename.endsWith('.gz'));
    await Promise.all(gz_files.map(async (filename) => {
      const input_path = `${input_dir}/${filename}`;
      const output_path = `${input_unzipped_dir}/${filename.slice(0, -3)}`;
      const readStream = fs.createReadStream(input_path);
      const unzip = zlib.createGunzip();
      const writeStream = fs.createWriteStream(output_path);
      await pipeline(readStream, unzip, writeStream);
    }));
    console.log('解压完成');

    // 步骤2:提取、转换、保存、压缩JSON文件
    const json_files = (await get_files(input_unzipped_dir)).filter(filename => filename.endsWith('.json'));
    await Promise.all(json_files.map(async (filename) => {
      const input_path = `${input_unzipped_dir}/${filename}`;
      const raw_data = await read_json_file(input_path);
      const transformed_data = transform_JSON(raw_data);
      const transformed_str = JSON.stringify(transformed_data, null, 2);

      // 保存转换后的文件
      const output_path = `${output_dir}/${filename}`;
      await writeFile(output_path, transformed_str, 'utf-8');

      // 压缩转换后的文件
      await gzip_file(output_path, `${output_path}.gz`);
      console.log(`处理完成:${filename}`);
    }));
    console.log('ETL管道执行完成');

  } catch (error) {
    console.error('管道执行出错:', error);
  }
};

orchestrate_etl_pipeline();

关键改进点

  • 用stream.promises.pipeline替代手动监听流事件,简化异步流操作的错误处理。
  • 全程使用async/await统一异步逻辑,确保步骤按顺序执行。
  • 将回调式函数改为Promise风格,避免回调地狱和异步等待问题。
  • 过滤非目标文件,避免无效Promise混入。
  • 保留转换函数的纯函数特性,避免变量覆盖。
  • 补充了保存转换结果和重新压缩的完整逻辑。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 20:46:35