NodeJS Stream Pipeline报错:val参数需为Readable/Iterable实例
问题场景
运行环境为 Node v16.15.0,执行index3.js脚本时,使用stream/promises模块导出的pipeline方法对接文件读流与写流,触发如下类型错误:
TypeError [ERR_INVALID_ARG_TYPE]: The "val" argument must be an instance of Readable, Iterable, or AsyncIterable. Received an instance of WriteStream
问题代码
文件:index3.js
#!/usr/bin/env node 'use strict'; import { pipeline } from 'stream/promises' import { realpathSync, createReadStream, createWriteStream } from 'fs'; import { pathToFileURL } from 'url'; async function doStuff() { return new Promise((resolve, reject) => { let readStream = createReadStream("input.js"); let writeStream = createWriteStream("output.js"); pipeline( readStream, writeStream, async(err) => { if (err) { console.error('failed', err); reject({res:'Pipeline failed', err}); } else { console.log('succeeded'); resolve('succeeded'); } } ); }); } export default function myFunc() { doStuff().catch(err => console.log(err)); } function wasCalledAsScript() { const realPath = realpathSync(process.argv[1]); const realPathAsUrl = pathToFileURL(realPath).href; return import.meta.url === realPathAsUrl; } if (wasCalledAsScript()) { myFunc(); }
运行命令与报错日志
# node index3.js node:internal/errors:465 ErrorCaptureStackTrace(err); ^ TypeError [ERR_INVALID_ARG_TYPE]: The "val" argument must be an instance of Readable, Iterable, or AsyncIterable. Received an instance of WriteStream at new NodeError (node:internal/errors:372:5) at makeAsyncIterable (node:internal/streams/pipeline:100:9) at pipelineImpl (node:internal/streams/pipeline:263:13) at node:stream/promises:28:5 at new Promise (<anonymous>) at pipeline (node:stream/promises:17:10) at file:///var/www/html/index3.js:14:9 at new Promise (<anonymous>) at doStuff (file:///var/www/html/index3.js:10:12) at myFunc (file:///var/www/html/index3.js:32:5) { code: 'ERR_INVALID_ARG_TYPE' }
报错根因
混淆了两个不同版本pipeline的调用签名:
- 回调版
pipeline(从stream模块直接导入)的签名是pipeline(...streams, callback),最后一个参数是执行完成的回调函数 - Promise版
pipeline(从stream/promises导入)的签名是pipeline(...streams),不接收回调参数,方法执行后直接返回Promise,所有传入的参数都会被识别为流链路节点
代码中给Promise版pipeline传了3个参数:读流、写流、回调函数。pipeline默认按链路规则校验参数:除最后一个参数必须是可写流外,前面所有参数都必须是可读流/可迭代对象/双工转换流。此时第二个参数WriteStream被识别为链路中间节点,它不属于可读类型,直接触发参数类型错误。
另外手动包裹一层Promise属于冗余逻辑,Promise版pipeline本身已经返回Promise,不需要重复包装。
正确实现
基础读写流对接示例
去掉回调参数,直接await pipeline的执行结果即可:
#!/usr/bin/env node 'use strict'; import { pipeline } from 'stream/promises' import { realpathSync, createReadStream, createWriteStream } from 'fs'; import { pathToFileURL } from 'url'; async function doStuff() { const readStream = createReadStream("input.js"); const writeStream = createWriteStream("output.js"); // 直接await,不需要传回调 await pipeline(readStream, writeStream); console.log('succeeded'); return 'succeeded'; } export default function myFunc() { doStuff().catch(err => { console.error('failed', err); }); } function wasCalledAsScript() { const realPath = realpathSync(process.argv[1]); const realPathAsUrl = pathToFileURL(realPath).href; return import.meta.url === realPathAsUrl; } if (wasCalledAsScript()) { myFunc(); }
带Transform流的逐行处理示例
如果需要在链路中加入内容转换逻辑,直接把Transform流插到读流和写流中间即可,以下是逐行给文本加行号的示例:
#!/usr/bin/env node 'use strict'; import { pipeline } from 'stream/promises' import { realpathSync, createReadStream, createWriteStream } from 'fs'; import { pathToFileURL } from 'url'; import { Transform } from 'stream'; import readline from 'readline/promises'; async function doStuff() { const readStream = createReadStream("input.js", { encoding: 'utf8' }); const writeStream = createWriteStream("output.js", { encoding: 'utf8' }); // 逐行读取的异步迭代器 const rl = readline.createInterface({ input: readStream, crlfDelay: Infinity }); // 转换流:逐行处理内容 const lineTransform = new Transform({ writableObjectMode: true, readableObjectMode: false, transform(line, _, callback) { // 示例处理:给每一行前面加行号 const processedLine = `// ${this.lineNumber++}: ${line}\n`; callback(null, processedLine); }, construct(callback) { this.lineNumber = 1; callback(); } }); // 链路顺序:读流 -> 逐行迭代 -> 转换流 -> 写流 await pipeline(rl, lineTransform, writeStream); console.log('处理完成'); return 'succeeded'; } export default function myFunc() { doStuff().catch(err => { console.error('处理失败', err); }); } function wasCalledAsScript() { const realPath = realpathSync(process.argv[1]); const realPathAsUrl = pathToFileURL(realPath).href; return import.meta.url === realPathAsUrl; } if (wasCalledAsScript()) { myFunc(); }
注意:pipeline会在执行完成/出错时自动销毁链路中所有流,不需要手动调用
destroy或监听close事件做资源清理。
内容的提问来源于stack exchange,提问作者Tyler V.
相关产品推荐
相关产品推荐

