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
相关产品推荐
相关产品推荐

