
本文详解如何使用 knex.js 高效、低内存地将 200 万行 csv 数据批量写入 sqlite,避免事件流中直接插入导致的并发错误与内存溢出问题,并提供可落地的分块批处理方案。
本文详解如何使用 knex.js 高效、低内存地将 200 万行 csv 数据批量写入 sqlite,避免事件流中直接插入导致的并发错误与内存溢出问题,并提供可落地的分块批处理方案。
在处理大规模 CSV 文件(如含 200 万行数据)时,直接在 fast-csv 的 'data' 事件回调中逐行调用 knex.insert() 是高危操作:它会瞬间触发大量异步数据库写入,既无法控制并发量,又持续累积未完成 Promise,极易引发 SQLite 的 SQLITE_BUSY 错误、连接池耗尽或 Node.js 内存暴涨(V8 堆溢出)。原代码虽将数据暂存于数组再统一 batchInsert,但一次性加载全部 200 万行对象到内存(尤其含字符串字段时),仍可能导致数百 MB 内存占用,缺乏可扩展性。
✅ 推荐方案:流式分块批处理(Streaming + Chunked Batch Insert)
核心原则是——不缓存全量数据,不并发乱序写入,而是在流中按固定大小切片、顺序执行可控批次插入。以下是优化后的健壮实现:
import * as fs from 'fs';
import * as fastcsv from 'fast-csv';
import { knex } from './knex-config'; // 确保已配置 SQLite 连接池(min: 1, max: 2)
async function insertCsvToDb(filePath: string, tableName: string, batchSize = 5000): Promise<void> {
const stream = fs.createReadStream(filePath);
const parser = fastcsv.parse({ headers: true, ignoreEmpty: true });
let chunk: any[] = [];
let totalInserted = 0;
// 使用 Promise 链确保严格串行执行,避免并发冲突
const insertChunk = async (): Promise<void> => {
if (chunk.length === 0) return;
try {
await knex(tableName).insert(chunk);
totalInserted += chunk.length;
console.log(`✓ Inserted batch of ${chunk.length} rows. Total: ${totalInserted}`);
chunk = []; // 清空当前块
} catch (err) {
console.error(`✗ Batch insert failed at row ${totalInserted}:`, err);
throw err;
}
};
// 流式处理:边读边攒批,不全量加载
for await (const row of parser.fromStream(stream)) {
chunk.push(row);
if (chunk.length >= batchSize) {
await insertChunk(); // 等待当前批完成,再处理下一批
}
}
// 插入剩余不足 batchSize 的尾部数据
await insertChunk();
console.log(`✅ All ${totalInserted} rows inserted successfully.`);
}</void></void>
? 关键优化点说明:
- 内存友好:for await...of 配合 fast-csv 的可迭代流接口,全程仅保留一个 batchSize(如 5000)大小的内存缓冲区,峰值内存 ≈ 单批数据体积(远低于 200 万行全载);
- 事务安全:每批次独立 insert(),利用 SQLite 的自动事务机制(单语句即事务),避免手动事务嵌套复杂度;若需更强一致性,可在外层包裹 knex.transaction();
- 错误可控:失败时精确报错位置(当前已插入总数),支持断点续插(需配合记录偏移量);
- 连接池适配:SQLite 轻量但不支持高并发写,将连接池 max 设为 1–2,配合串行批处理,彻底规避 SQLITE_BUSY;
⚠️ 注意事项:
- 避免在 .on('data') 中直接 await 或 .then() —— 事件回调非 Promise-aware,易造成“幽灵”未处理 Promise;
- batchSize 需权衡:过大(>10000)可能触发 SQLite 单语句参数上限(SQLITE_MAX_VARIABLE_NUMBER,默认 999);过小(
- 对超大字段(如长文本、JSON),务必启用 parser.on('error', ...) 全局捕获解析异常,防止流中断;
- 生产环境建议添加进度条(如 progress 库)与写入速率统计,便于监控。
通过该方案,200 万行 CSV 可在数分钟内稳定导入,内存占用稳定在 50–100MB 区间(取决于行宽),真正实现高性能、低开销、生产就绪的数据迁移。











