1. ストリームの種類
| 種類 | 説明 | たとえば |
|---|---|---|
| 読みやすい | データソース | fs.createReadStream、http リクエスト |
| 書き込み可能 | データ記録先 | fs.createWriteStream、http 応答 |
| トランスフォーム | 読み取り + 変換 + 書き込み | zlib.createGzip、crypto.createCipher |
| デュプレックス | 読み取りと書き込みを独立して行う | net.Socket、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。