You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.16 22:53:12