返回博客
技术 2025年3月18日 8 分钟阅读 · 1799 字

Node.js Stream 流处理与 Buffer 操作详解

掌握流式数据处理核心,构建高效的大文件与数据管道

#Node.js #Stream #Buffer #流处理
本文由 AI 辅助生成,经人工审核发布

在处理大文件或高频数据流时,一次性将所有数据加载到内存会导致严重的性能问题。Node.js 提供了 Stream(流)和 Buffer(缓冲区)两大核心模块,让我们可以分块处理数据,实现高效的内存管理和数据管道。本文将深入讲解两者的用法,并通过实战案例掌握流式编程。

一、Buffer:二进制数据容器

Buffer 是 Node.js 中处理二进制数据的基础,它在 V8 堆内存之外分配,专门用于存储原始字节序列。

1.1 创建 Buffer

// 从字符串创建(默认 UTF-8 编码)
const buf1 = Buffer.from('Hello 世界', 'utf8');
console.log(buf1); // <Buffer 48 65 6c 6c 6f 20 e4 b8 96 e7 95 8c>
console.log(buf1.length); // 12(字节数,不是字符数)

// 从数组创建
const buf2 = Buffer.from([0x48, 0x65, 0x6c, 0x6c, 0x6f]);
console.log(buf2.toString()); // Hello

// 分配指定大小的空 Buffer
const buf3 = Buffer.alloc(10); // 填充 0
const buf4 = Buffer.allocUnsafe(10); // 不初始化,更快但有残留数据

注意Buffer.allocUnsafe 不会清零内存,速度更快但可能包含敏感数据,仅在确定会完全覆写时使用。

1.2 读写操作

const buf = Buffer.alloc(16);

// 写入数据
buf.write('Hello', 0, 'utf8');        // 从偏移 0 写入字符串
buf.writeUInt8(255, 5);                // 在偏移 5 写入单字节
buf.writeUInt16BE(1000, 6);            // 大端序 16 位整数
buf.writeDoubleBE(Math.PI, 8);         // 双精度浮点数

// 读取数据
console.log(buf.toString('utf8', 0, 5));  // Hello
console.log(buf.readUInt8(5));             // 255
console.log(buf.readUInt16BE(6));          // 1000
console.log(buf.readDoubleBE(8));          // 3.141592653589793

1.3 常用操作

const a = Buffer.from('Hello ');
const b = Buffer.from('World');

// 拼接
const combined = Buffer.concat([a, b]);
console.log(combined.toString()); // Hello World

// 比较
console.log(Buffer.compare(a, b)); // -1(a < b)

// 切片(共享内存)
const sliced = combined.subarray(0, 5);
console.log(sliced.toString()); // Hello

// 复制
const copied = Buffer.alloc(5);
combined.copy(copied, 0, 0, 5);

Buffer 的常见场景:文件 I/O、网络传输、加密计算。以下是一个 Base64 编解码的示例:

const text = 'Hello 世界';

// 编码
const encoded = Buffer.from(text, 'utf8').toString('base64');
console.log(encoded); // SGVsbG8g5LiW55+l

// 解码
const decoded = Buffer.from(encoded, 'base64').toString('utf8');
console.log(decoded); // Hello 世界

二、四种流类型

Node.js 的 Stream 模块定义了四种流类型:

类型描述典型场景
Readable可读流文件读取、HTTP 请求体
Writable可写流文件写入、HTTP 响应体
Duplex双工流TCP Socket
Transform转换流Gzip 压缩、加密
const { Readable, Writable, Duplex, Transform } = require('stream');

// Readable 示例
const readable = Readable.from(['Hello', ' ', 'World']);

// Writable 示例
const writable = new Writable({
  write(chunk, encoding, callback) {
    console.log('写入:', chunk.toString());
    callback();
  }
});

readable.pipe(writable);
// 输出: 写入: Hello  写入: World

三、pipe 管道链

pipe 方法将可读流连接到可写流,自动处理数据流动和背压。多个流可以通过管道串联,形成处理链。

const fs = require('fs');
const zlib = require('zlib');
const crypto = require('crypto');

// 读取文件 → Gzip 压缩 → 加密 → 写入文件
fs.createReadStream('input.txt')
  .pipe(zlib.createGzip())                       // 压缩
  .pipe(crypto.createCipheriv('aes-256-cbc', key, iv)) // 加密
  .pipe(fs.createWriteStream('output.enc'));     // 写入

// 更推荐使用 pipeline(更好的错误处理)
const { pipeline } = require('stream/promises');

async function compressAndEncrypt() {
  await pipeline(
    fs.createReadStream('input.txt'),
    zlib.createGzip(),
    fs.createWriteStream('output.gz')
  );
  console.log('处理完成');
}

pipeline 相比 pipe 的优势:

特性pipepipeline
错误处理手动监听 error自动传递错误
资源清理需手动销毁自动销毁所有流
代码风格链式函数式
Promise 支持

四、自定义流实现

4.1 自定义 Readable

const { Readable } = require('stream');

class CounterStream extends Readable {
  constructor(max, options) {
    super(options);
    this.max = max;
    this.current = 0;
  }

  _read() {
    if (this.current < this.max) {
      const chunk = Buffer.from(`${this.current}\n`, 'utf8');
      this.push(chunk);
      this.current++;
    } else {
      this.push(null); // 结束流
    }
  }
}

const counter = new CounterStream(10);
counter.on('data', data => process.stdout.write(data));
counter.on('end', () => console.log('计数完成'));

4.2 自定义 Transform

Transform 流是 Duplex 的特化,它接收输入、变换后输出。

const { Transform } = require('stream');

// 行转大写
class UpperCaseStream extends Transform {
  _transform(chunk, encoding, callback) {
    this.push(chunk.toString().toUpperCase());
    callback();
  }
}

// JSON 行解析器
class JsonLineParser extends Transform {
  constructor() {
    super({ objectMode: true });
    this.buffer = '';
  }

  _transform(chunk, encoding, callback) {
    this.buffer += chunk.toString();
    const lines = this.buffer.split('\n');
    this.buffer = lines.pop(); // 保留不完整的最后一行

    for (const line of lines) {
      if (line.trim()) {
        try {
          this.push(JSON.parse(line));
        } catch {
          callback(new Error(`无效的 JSON: ${line}`));
          return;
        }
      }
    }
    callback();
  }

  _flush(callback) {
    if (this.buffer.trim()) {
      this.push(JSON.parse(this.buffer));
    }
    callback();
  }
}

// 使用:日志文件解析
const logStream = fs.createReadStream('app.log');
logStream
  .pipe(new JsonLineParser())
  .on('data', obj => console.log('日志条目:', obj.level, obj.message));

五、背压处理

当数据生产速度超过消费速度时,会产生背压。如果忽视背压,数据会在内存中堆积,最终导致内存溢出。

// 错误示范:忽略背压
readable.on('data', chunk => {
  writable.write(chunk); // 如果写入慢,数据会堆积
});

// 正确示范:监听 drain 事件
readable.on('data', chunk => {
  const ok = writable.write(chunk);
  if (!ok) {
    readable.pause();
    writable.once('drain', () => readable.resume());
  }
});

// 推荐方式:使用 pipe 或 pipeline 自动处理背压
readable.pipe(writable);

write() 方法返回 false 表示内部缓冲已满,此时应暂停读取,等待 drain 事件后再继续。pipe 内部已实现这一机制。

六、大文件处理实战

6.1 逐行读取

const { createReadStream } = require('fs');
const { createInterface } = require('readline');

async function processLargeFile(filePath) {
  const fileStream = createReadStream(filePath);
  const rl = createInterface({
    input: fileStream,
    crlfDelay: Infinity
  });

  let lineCount = 0;
  let totalBytes = 0;

  for await (const line of rl) {
    lineCount++;
    totalBytes += Buffer.byteLength(line);
    if (lineCount % 100000 === 0) {
      console.log(`已处理 ${lineCount} 行`);
    }
  }

  console.log(`完成: 共 ${lineCount} 行, ${totalBytes} 字节`);
}

processLargeFile('huge-data.csv');

6.2 流式 CSV 解析

const { pipeline, Transform } = require('stream');
const { parse } = require('csv-parse');

async function analyzeCSV(inputPath, outputPath) {
  let rowCount = 0;
  let totalAmount = 0;

  const filter = new Transform({
    objectMode: true,
    transform(row, encoding, callback) {
      const amount = Number(row[2]);
      if (amount > 1000) {
        rowCount++;
        totalAmount += amount;
        this.push(row);
      }
      callback();
    }
  });

  await pipeline(
    fs.createReadStream(inputPath),
    parse({ columns: true, delimiter: ',' }),
    filter,
    fs.createWriteStream(outputPath)
  );

  console.log(`筛选 ${rowCount} 条记录,总金额 ${totalAmount}`);
}

6.3 HTTP 流式上传

const http = require('http');
const fs = require('fs');

const server = http.createServer((req, res) => {
  if (req.method === 'POST' && req.url === '/upload') {
    // req 本身就是 Readable 流
    const writeStream = fs.createWriteStream('./uploads/uploaded.dat');

    req.pipe(writeStream);

    writeStream.on('finish', () => {
      res.writeHead(200, { 'Content-Type': 'application/json' });
      res.end(JSON.stringify({ message: '上传成功' }));
    });

    writeStream.on('error', err => {
      res.writeHead(500);
      res.end('上传失败');
    });
  }
});

七、HTTP 流式响应

流式响应可以实现服务端推送(Server-Sent Events)或分块传输,避免长时间等待。

const http = require('http');

http.createServer((req, res) => {
  if (req.url === '/stream') {
    res.writeHead(200, {
      'Content-Type': 'text/event-stream',
      'Cache-Control': 'no-cache',
      'Connection': 'keep-alive'
    });

    let id = 0;
    const interval = setInterval(() => {
      const data = { time: new Date(), id: id++ };
      res.write(`data: ${JSON.stringify(data)}\n\n`);

      if (id > 10) {
        res.write('event: close\n');
        res.write('data: 结束\n\n');
        res.end();
        clearInterval(interval);
      }
    }, 1000);

    req.on('close', () => clearInterval(interval));
  }
}).listen(3000);

大文件下载示例:

const fileStream = fs.createReadStream('large-video.mp4');

res.writeHead(200, {
  'Content-Type': 'video/mp4',
  'Content-Disposition': 'attachment; filename="video.mp4"'
});

fileStream.pipe(res);

// 支持断点续传
const range = req.headers.range;
if (range) {
  const [start, end] = range.replace(/bytes=/, '').split('-').map(Number);
  fileStream.destroy();
  const partial = fs.createReadStream('large-video.mp4', { start, end });
  res.writeHead(206, {
    'Content-Range': `bytes ${start}-${end}/*`,
    'Accept-Ranges': 'bytes'
  });
  partial.pipe(res);
}

八、objectMode 对象流

默认流处理的是 Buffer 或字符串。设置 objectMode: true 后,流可以处理任意 JavaScript 对象。

const { Readable } = require('stream');

const users = [
  { id: 1, name: 'Alice' },
  { id: 2, name: 'Bob' },
  { id: 3, name: 'Charlie' }
];

const userStream = Readable.from(users); // 默认 objectMode

userStream.on('data', user => {
  console.log(`用户: ${user.name}`); // 直接是对象,无需 JSON.parse
});

总结

Stream 和 Buffer 是 Node.js 处理 I/O 密集型任务的两大基石。Buffer 解决了二进制数据的存储与转换问题,Stream 则提供了分块处理数据的流式模型。掌握管道链、自定义流、背压处理和 objectMode,就能构建出高效的大文件处理管道和高性能的 HTTP 流式服务。在实际开发中,优先使用 pipeline 替代 pipe,利用 stream/promises 的 Promise 接口简化异步流程,能让代码更健壮、更易维护。