Chuyển đến nội dung chính

Lesson 6: Streams & Buffers

Readable, Writable, Transform, Duplex streams. Backpressure, pipeline(), stream.compose(). Buffer API, ArrayBuffer, TypedArrays. Stream-based file processing, CSV/JSON parsing.

💻 Programming — Lesson 6 Lesson 6: Streams & Buffers

Node.js Core: From Basics to Advanced

Part 2: Core Modules Deep Dive

xdev.asia

1. Types of Streams

TypeDescriptionFor example
ReadableData sourcefs.createReadStream, http request
WriteableData recording destinationfs.createWriteStream, http response
TransformRead + transform + writezlib.createGzip, crypto.createCipher
DuplexRead + write independentlynet.Socket, WebSocket

2. Custom Readable Stream

import { Readable } from 'node:stream'

// Readable từ array
const readable = Readable.from(['Hello', ' ', 'World'])

for await (const chunk of readable) {
  console.log(chunk) // 'Hello', ' ', 'World'
}

// Custom Readable
class CounterStream extends Readable {
  private current = 0
  constructor(private max: number) {
    super({ objectMode: true })
  }
  _read() {
    if (this.current < this.max) {
      this.push({ value: this.current++ })
    } else {
      this.push(null) // Signal end
    }
  }
}

3. Transform Stream

import { Transform, pipeline } from 'node:stream'
import { createReadStream, createWriteStream } from 'node:fs'
import { promisify } from 'node:util'

const pipelineAsync = promisify(pipeline)

// Transform: uppercase mỗi dòng
const toUpperCase = new Transform({
  transform(chunk, encoding, callback) {
    callback(null, chunk.toString().toUpperCase())
  },
})

// CSV line parser
const csvParser = new Transform({
  objectMode: true,
  transform(chunk, encoding, callback) {
    const lines = chunk.toString().split('\n').filter(Boolean)
    for (const line of lines) {
      const [name, age, city] = line.split(',')
      this.push({ name, age: Number(age), city })
    }
    callback()
  },
})

await pipelineAsync(
  createReadStream('input.txt'),
  toUpperCase,
  createWriteStream('output.txt')
)

4. Backpressure

import { createReadStream, createWriteStream } from 'node:fs'

const readable = createReadStream('large-file.bin')
const writable = createWriteStream('output.bin')

// Xử lý backpressure thủ công
readable.on('data', (chunk) => {
  const canContinue = writable.write(chunk)
  if (!canContinue) {
    readable.pause() // Tạm dừng đọc
    writable.once('drain', () => readable.resume()) // Resume khi đã ghi xong
  }
})

// pipeline() tự xử lý backpressure — KHUYẾN NGHỊ
import { pipeline } from 'node:stream/promises'
await pipeline(readable, writable)

5. Buffer API

// Tạo Buffer
const buf1 = Buffer.from('Hello', 'utf8')
const buf2 = Buffer.alloc(16)         // Zero-filled
const buf3 = Buffer.allocUnsafe(16)   // Không clear — nhanh hơn

// Đọc/ghi
buf2.writeUInt32BE(0x48454c4c, 0) // 'HELL'
console.log(buf1.toString('hex'))   // 48656c6c6f
console.log(buf1.toString('base64'))

// So sánh
Buffer.compare(buf1, buf2) // -1, 0, 1
buf1.equals(buf2)           // false

// Nối
const combined = Buffer.concat([buf1, buf2])

6. stream.compose (Node.js 22+)

import { compose } from 'node:stream'

// Compose nhiều transforms thành 1 stream
const processPipeline = compose(
  csvParser,
  filterAge,
  toJSON
)

await pipeline(
  createReadStream('users.csv'),
  processPipeline,
  createWriteStream('output.json')
)

Next article: HTTP/HTTPS & HTTP/2 — server, routing, TLS.