1. 流的類型
| 類型 | 描述 | 例如 |
|---|---|---|
| 可讀 | 數據來源 | fs.createReadStream,http請求 |
| 可寫 | 數據記錄目的地 | fs.createWriteStream,http 回應 |
| 變換 | 讀取+轉換+寫入 | zlib.createGzip、crypto.createCipher |
| 複式 | 獨立讀+寫 | 網路套接字、WebSocket |
2. 自訂可讀串流
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. 變換流
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. 背壓
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. 緩衝區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')
)
下一篇: HTTP/HTTPS 和 HTTP/2 — 伺服器、路由、TLS。