Không đọc cả file vào RAM.
Dựng một pipeline đọc → parse → gom lô → ghi, để backpressure tự điều tiết tốc độ theo DB.
import { pipeline } from 'node:stream/promises'
import { createReadStream } from 'node:fs'
import { Transform, Writable } from 'node:stream'
import { parse } from 'csv-parse'
const BATCH = 1000
const batcher = new Transform({
objectMode: true,
transform(row, _enc, cb) {
this.rows ??= []
this.rows.push(row)
if (this.rows.length >= BATCH) { const b = this.rows; this.rows = []; return cb(null, b) }
cb()
},
flush(cb) { cb(null, this.rows?.length ? this.rows : undefined) },
})
const writer = new Writable({
objectMode: true,
async write(batch, _enc, cb) {
try { await db.insertMany(batch); cb() } catch (err) { cb(err) }
},
})
await pipeline(createReadStream('data.csv'), parse({ columns: true }), batcher, writer)Vì sao không hết bộ nhớ: writer._write chỉ gọi cb() sau khi DB ghi xong. Chừng nào chưa gọi, hàng đợi nội bộ đầy tới highWaterMark, write() trả false và stream phía trên ngừng đọc cho tới khi có drain. Bộ nhớ giữ ở mức một lô, không phụ thuộc kích thước file.
Các điểm cần nói thêm khi phỏng vấn:
- Ghi theo lô (500–2000 dòng) thay vì từng dòng để giảm round-trip.
- Idempotency: dùng ON CONFLICT DO NOTHING hoặc unique key để chạy lại an toàn khi job hỏng giữa chừng.
- Ghi checkpoint số dòng đã xử lý để có thể tiếp tục.
- Nếu import chạy trong API service, đẩy sang worker/queue để không chiếm event loop của các request khác.