如何在Node.js中构建完整ETL管道?解决执行顺序与JSON解析问题
Node.js ETL管道脚本修复方案
问题根源
- 异步流程无序:解压操作未等待完成就执行后续步骤,导致读取到不完整的JSON文件,触发
SyntaxError: Unexpected end of JSON input。 - 回调与Promise兼容问题:回调式的文件读取未包装为Promise,无法被
Promise.all正确等待。 - 变量意外覆盖:转换步骤中错误地将
transform_JSON函数覆盖为返回的字符串。 - 无效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
相关产品推荐
相关产品推荐

