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

Node.js中如何统计Duplex流的读写字节数?

问题:统计接入HTTP请求的Duplex流总读写字节数

我有一个来自第三方库的Duplex流对象,要将其接入HTTP请求,请求与响应都通过该流传输。根据Node.js官方文档要求:开发者应仅选择一种方式消费单个流的数据,同时使用on('data')、on('readable')、pipe()或异步迭代器会导致异常行为。

现有代码示例如下,需要实现统计该Duplex流的总读写字节数:

const req = http.request({
    createConnection: () => {
        const duplexStream = getDuplexStreamFromLibrary();
        // TODO: Count read and written bytes in `duplexStream`
        return duplexStream;
    }
});

req.on("response", (res) => {
    res.pipe(anotherStream);
});

解决方案:代理Duplex流实现字节统计

因为不能直接在原流上添加data事件监听(会和HTTP模块的流消费方式冲突),所以可以创建一个代理Duplex流,用来包裹原第三方流,在代理层完成字节统计,同时不干扰原流的正常工作。

实现步骤

  • 继承stream.Duplex创建自定义代理流
  • 在代理流的_write方法中统计写入字节数,并将数据转发给原流
  • 在代理流的_read方法中从原流读取数据,统计读取字节数后再推送给下游
  • 转发原流的所有事件到代理流,保证流的行为一致

完整代码示例

const http = require('http');
const { Duplex } = require('stream');

const req = http.request({
    createConnection: () => {
        const duplexStream = getDuplexStreamFromLibrary();
        // 初始化统计变量
        let bytesWritten = 0;
        let bytesRead = 0;

        // 创建代理Duplex流
        const proxyStream = new Duplex({
            // 同步原流的核心配置,保证行为一致
            readableObjectMode: duplexStream.readableObjectMode,
            writableObjectMode: duplexStream.writableObjectMode,
            highWaterMark: duplexStream.highWaterMark,

            _write(chunk, encoding, callback) {
                // 统计写入字节:区分Buffer和字符串场景
                if (Buffer.isBuffer(chunk)) {
                    bytesWritten += chunk.length;
                } else {
                    bytesWritten += Buffer.byteLength(chunk, encoding);
                }
                // 将数据转发给原流
                duplexStream.write(chunk, encoding, callback);
            },

            _read(size) {
                // 从原流读取数据
                const chunk = duplexStream.read(size);
                if (chunk !== null) {
                    // 统计读取字节
                    bytesRead += Buffer.isBuffer(chunk) ? chunk.length : Buffer.byteLength(chunk);
                    // 推送给下游消费
                    this.push(chunk);
                }
            }
        });

        // 转发原流的所有关键事件,避免状态不一致
        duplexStream.on('error', (err) => proxyStream.emit('error', err));
        duplexStream.on('end', () => proxyStream.push(null));
        duplexStream.on('close', () => proxyStream.emit('close'));
        duplexStream.on('drain', () => proxyStream.emit('drain'));

        // 流关闭时输出统计结果,也可按需存储到其他位置
        proxyStream.on('close', () => {
            console.log(`总写入字节数: ${bytesWritten}`);
            console.log(`总读取字节数: ${bytesRead}`);
        });

        return proxyStream;
    }
});

req.on("response", (res) => {
    res.pipe(anotherStream);
});

关键说明

  • 代理流同步原流的objectMode、highWaterMark等配置,确保流的性能和行为与原流完全一致
  • 字节统计逻辑覆盖了Buffer和字符串两种数据类型,保证统计精度
  • 所有原流的事件都被转发到代理流,避免出现流状态异常、数据丢失等问题

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 11:31:08