Node.js处理百万JSON数据调用外部API入库的最优最快方案咨询
百万级数据API调用+入库最优方案
一、现有代码的核心瓶颈
你的当前实现存在几个致命效率问题:
- 一次性将百万条数据加载到内存,极易触发内存溢出(OOM)
- 用
await在for循环里调用API,本质是串行执行,完全浪费了Node.js的异步并发能力 - 数据库逐条插入,重复建立连接、执行SQL,IO开销极大
二、优化方案与代码实现
1. 大文件流式读取:避免内存爆炸
不用一次性读取整个JSON文件,改用流式读取+流式解析,边读边处理单条数据,内存占用始终维持在低水平。依赖JSONStream库实现:
安装依赖:
npm install jsonstream
示例代码:
const fs = require('fs'); const JSONStream = require('JSONStream'); // 流式读取JSON数组的每一项 function readJsonStream() { return fs.createReadStream('./datas.json', 'utf-8') .pipe(JSONStream.parse('*')); }
2. API调用:控制并发量,真正并行
用p-limit实现并发池,限制同时发起的API请求数(比如30个,根据目标API的限流规则调整),既保证并行效率,又不会触发API限流:
安装依赖:
npm install p-limit
示例代码:
const pLimit = require('p-limit'); const pRetry = require('p-retry'); // 可选:失败自动重试 // 限制同时30个API请求 const limit = pLimit(30); async function processApiCalls() { const token = await fetchToken(); const stream = readJsonStream(); const promises = []; let processedCount = 0; stream.on('data', (item) => { processedCount++; if (processedCount % 1000 === 0) { console.log(`已处理 ${processedCount} 条请求`); } // 包装受并发限制的API调用,加失败重试 const promise = limit(async () => { try { return await pRetry(async () => { const result = await siamOrganization(token, { identifier: { issuer: item.issuer, assigner: item.assigner, id: item.id, type: item.type, }, }); if (result.status.statusCode === 200 && result.data.data) { return result.data.data; } return null; }, { retries: 3 }); // 失败重试3次 } catch (err) { console.error(`请求失败(ID: ${item.id}):`, err.message); return null; } }); promises.push(promise); }); // 等待流式读取完成 await new Promise(resolve => stream.on('end', resolve)); // 过滤无效结果,返回有效数据 return (await Promise.all(promises)).filter(Boolean); }
3. 数据库入库:批量插入,减少IO开销
把数据攒成批量(比如1000条)再插入数据库,大幅减少连接和SQL执行次数,提升写入效率:
示例代码(以SQL Server的mssql库为例):
const sql = require('mssql'); async function batchInsert(results) { const batchSize = 1000; // 建立数据库连接池(只初始化一次) const pool = await sql.connect({ user: '你的用户名', password: '你的密码', server: '你的服务器地址', database: '你的数据库名', pool: { max: 10, min: 2, idleTimeoutMillis: 30000 } }); await createTableIfNotExists(pool); // 复用你的建表逻辑,传入连接池 for (let i = 0; i < results.length; i += batchSize) { const batch = results.slice(i, i + batchSize); try { const request = pool.request(); // 构建批量插入SQL(根据你的表结构调整字段) const valueStrings = batch.map(item => `('${item.issuer}', '${item.assigner}', '${item.id}', '${item.type}')` ).join(','); await request.query(` INSERT INTO 你的表名 (issuer, assigner, id, type) VALUES ${valueStrings} `); console.log(`已插入 ${i + batch.length} 条数据`); } catch (err) { console.error(`批量插入失败(批次:${i/batchSize +1}):`, err.message); // 可选:记录失败批次到文件,后续补插 } } await pool.close(); }
4. 其他细节优化
- Token缓存:如果token有效期较长,缓存起来,过期后再重新获取,避免重复调用
fetchToken - 进度监控:添加计数器,实时打印处理进度,方便跟踪
- 错误归档:将失败的请求和入库错误写入日志文件,后续统一排查补处理
整合完整流程
async function main() { try { console.log('开始处理数据...'); const validResults = await processApiCalls(); console.log(`API调用完成,共获取 ${validResults.length} 条有效数据`); await batchInsert(validResults); console.log('所有数据处理完成!'); } catch (err) { console.error('全局错误:', err); } } main();
内容的提问来源于stack exchange,提问作者Kasra Roshan
相关产品推荐
相关产品推荐

