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

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.

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.03 01:36:47