如何用Sails.js+Waterline ORM批量插入10万条PostgreSQL数据?
优化Sails.js+Waterline批量插入PostgreSQL的方案
现有代码的核心问题
- 全量加载文件到内存:
readFileSync一次性读取10万条数据,占用大量内存,引发频繁GC拖慢处理速度。 - 异步操作未等待:
Link.createEach(chunkArray).usingConnection(db)未加await,会导致并发插入触发连接池耗尽或外键约束错误。 - 分块逻辑不合理:循环依赖
users.length,若用户与链接数据量不一致会遗漏数据;1万条的块大小对PostgreSQL过大,易引发锁表和超时。 - 关联处理冗余:数据同时包含用户
links数组和链接user字段,重复解析会增加Waterline的关联处理开销。
优化方案
1. 流处理文件,逐行读取
用Node.js原生fs.createReadStream和readline模块边读边处理,避免一次性加载全量数据到内存,降低内存占用。
2. 优化批量插入策略
- 调整分块大小为1000条:PostgreSQL批量插入的最优区间,平衡插入速度与数据库负载
- 关闭不必要的模型生命周期钩子:如
beforeCreate/afterCreate,减少额外逻辑开销 - 使用
fetch: false参数:不需要返回插入后的记录,大幅提升插入速度 - 保证插入顺序:先插入用户(链接依赖用户外键),再插入链接
3. 事务与并发控制
- 将每个分块的插入放在独立事务中,避免单个大事务长时间占用数据库锁
- 添加短暂延迟:分块插入后暂停100ms,避免数据库压力过载
4. 关联数据处理
由于用户与链接一一对应,先批量插入用户并记录email与数据库实际id的映射,再插入链接时替换正确的外键值,避免外键约束错误。
完整优化代码示例
const fs = require('fs'); const readline = require('readline'); generateUsers: async (req, res) => { // 临时关闭模型生命周期钩子以提升速度 User.disableLifecycleCallbacks(); Link.disableLifecycleCallbacks(); const CHUNK_SIZE = 1000; let userCount = 0; let linkCount = 0; const userIdMap = new Map(); // 存储用户email -> 数据库实际id // 批量插入用户 const userStream = readline.createInterface({ input: fs.createReadStream('./api/db/users.txt'), crlfDelay: Infinity }); let userChunk = []; for await (const line of userStream) { if (!line.trim()) continue; // 跳过空行 const userData = JSON.parse(line); delete userData.links; // 移除关联字段,后续手动处理链接 userChunk.push(userData); if (userChunk.length >= CHUNK_SIZE) { const insertedUsers = await User.createEach(userChunk).fetch(); insertedUsers.forEach(user => userIdMap.set(user.email, user.id)); userCount += insertedUsers.length; userChunk = []; await new Promise(resolve => setTimeout(resolve, 100)); // 缓解数据库压力 } } // 处理剩余用户数据 if (userChunk.length > 0) { const insertedUsers = await User.createEach(userChunk).fetch(); insertedUsers.forEach(user => userIdMap.set(user.email, user.id)); userCount += insertedUsers.length; } // 批量插入链接 const linkStream = readline.createInterface({ input: fs.createReadStream('./api/db/links.txt'), crlfDelay: Infinity }); let linkChunk = []; for await (const line of linkStream) { if (!line.trim()) continue; const linkData = JSON.parse(line); // 替换为数据库实际用户id(根据你的数据结构调整映射逻辑) linkData.user = userIdMap.get(linkData.user); linkChunk.push(linkData); if (linkChunk.length >= CHUNK_SIZE) { await Link.createEach(linkChunk).fetch(false); // fetch: false不返回结果,提升速度 linkCount += CHUNK_SIZE; linkChunk = []; await new Promise(resolve => setTimeout(resolve, 100)); } } // 处理剩余链接数据 if (linkChunk.length > 0) { await Link.createEach(linkChunk).fetch(false); linkCount += linkChunk.length; } // 恢复模型生命周期钩子 User.enableLifecycleCallbacks(); Link.enableLifecycleCallbacks(); return res.ok(`${userCount} users and ${linkCount} links were generated and inserted.`); }
额外优化建议
- 调整PostgreSQL配置:增大
max_wal_size、work_mem参数,提升批量写入性能 - 临时禁用唯一约束:插入前临时关闭
email/phone/url的唯一约束,插入完成后重新启用(需确保数据无重复) - 使用原生SQL插入:若Waterline的
createEach仍不够快,可直接用sails.sendNativeQuery执行原生INSERT语句,进一步提升性能
内容的提问来源于stack exchange,提问作者Yurii Pakhota
相关产品推荐
相关产品推荐

