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

Node.js实现JobID批量处理与数据库分批写入方案问询

实现Node.js分批写入数据库的解决方案

看起来你需要解决的核心问题是异步请求下的批量数据收集与分批写入,原来的代码里因为https.request是异步操作,直接用counter判断分批会因为回调执行顺序的问题导致逻辑失效,而且没有正确收集批量数据。下面我会一步步帮你修改代码,实现每100条(或最后剩余的)数据批量写入数据库的逻辑:

第一步:封装异步请求为Promise

首先把https.request封装成Promise,这样可以用async/await来控制流程,避免回调地狱,也能确保我们在拿到数据后再进行后续的收集和判断:

function fetchJobStatus(jobId) {
  return new Promise((resolve, reject) => {
    const formJobStatusURL = `https://localhost:8091/api/job/${jobId}/status`;
    const option = {
      method: 'GET',
      headers: {
        'X-Message-Created-Ts': `${new Date().toISOString()}`,
        'X-Transaction-Created-Ts': `${new Date().toISOString()}`,
        'X-User-Id': 'PerformanceExecuter',
        'X-Client-Id': `${uuid()}`,
        'X-Message-Id': `${uuid()}`,
        'X-Transaction-Id': `${uuid()}`,
        'Content-Type': 'application/json'
      }
    };

    let content = '';
    const reqGet = https.request(formJobStatusURL, option, (response) => {
      response.on('data', (data) => {
        content += data;
      });
      response.on('end', () => {
        try {
          const jsonPayload = JSON.parse(content);
          // 提取需要的字段,返回metrics对象,用可选链避免字段不存在报错
          resolve({
            job_id: jsonPayload.id,
            status: jsonPayload.status,
            initiatedBy: jsonPayload.initiatedBy,
            product: jsonPayload.product,
            operation: jsonPayload.operation,
            startTimestamp: jsonPayload.startTimestamp,
            endTimestamp: jsonPayload.endTimestamp,
            totalRecords: jsonPayload.file?.totalRecords || null,
            failedRecords: jsonPayload.file?.totalFailedRecords || null,
            sucessRecords: jsonPayload.file?.totalSuccessRecords || null,
            inprogressRecords: jsonPayload.file?.totalInProgressRecords || null,
            sucessStatus: jsonPayload.results?.successFileAvailable || null,
            failureStatus: jsonPayload.results?.failureFileAvailable || null,
            uploadJob: jsonPayload.actionsAvailable?.dataUploadAllowed || null,
            abortJob: jsonPayload.actionsAvailable?.abortJob || null
          });
        } catch (err) {
          reject(new Error(`解析Job ${jobId}响应失败: ${err.message}`));
        }
      });
    });

    reqGet.on('error', (err) => {
      reject(new Error(`请求Job ${jobId}状态失败: ${err.message}`));
    });

    reqGet.end();
  });
}

第二步:实现分批收集与写入逻辑

接下来修改主函数,创建一个数组来批量收集metrics,每累计100条就写入数据库,循环结束后还要处理最后一批不足100条的数据:

const fileStatusprocess = require('../controller/readResultFile');
const config = require('../config/config');
const https = require('https');
const uuid = require('uuid-random');

// 先添加上面的fetchJobStatus函数

async function jobProcesser() {
  try {
    const iterator = await fileStatusprocess.processResultFile("C:\\Support\\result.csv"); // 注意转义反斜杠
    const jobIds = iterator[0];
    console.log('Total jobs are', jobIds.length);

    const batchMetrics = [];
    const batchSize = 100; // 每批100条

    for (let i = 0; i < jobIds.length; i++) {
      const jobId = jobIds[i];
      try {
        // 等待请求完成,拿到当前job的metrics
        const metrics = await fetchJobStatus(jobId);
        batchMetrics.push(metrics);

        // 判断是否达到批量大小,或者是最后一条数据
        if (batchMetrics.length === batchSize || i === jobIds.length - 1) {
          console.log(`Writing ${batchMetrics.length} records to database`);
          // 这里替换成你的InfluxDB写入逻辑,比如调用你的写入函数
          // await writeToInfluxDB(batchMetrics);
          
          // 写入完成后清空数组
          batchMetrics.length = 0;
        }
      } catch (err) {
        // 处理单个job请求失败的情况,避免整个流程中断
        console.error(`处理Job ${jobId}失败: ${err.message}`);
        // 可选:可以把失败的job记录下来后续重试
      }
    }

    console.log('所有Job处理完成');
  } catch (err) {
    console.error('读取文件失败:', err.message);
  }
}

jobProcesser();

关键修改点说明

  • Promise封装请求:把异步的https.request转为Promise,用await确保我们拿到数据后再进行收集,避免了原来回调中counter判断失效的问题。
  • 批量收集数组:用batchMetrics数组来存储每批的metrics,而不是每次创建单个对象。
  • 分批写入判断:每次添加数据后,检查数组长度是否达到100,或者是否是最后一个job(处理剩余不足100的情况),满足条件就执行写入,然后清空数组。
  • 错误处理:添加了请求失败、解析失败的错误捕获,避免单个job的问题导致整个流程崩溃,你还可以扩展失败job的重试逻辑。
  • 路径转义:注意文件路径的反斜杠要转义("C:\\Support\\result.csv"),否则会被当成转义字符处理。

可选优化:并发请求

如果你的服务支持并发请求,为了提高处理速度,可以用Promise.all来批量请求,比如每次请求100个job,然后一起写入数据库,这样效率会更高,示例代码如下(可选):

async function jobProcesser() {
  try {
    const iterator = await fileStatusprocess.processResultFile("C:\\Support\\result.csv");
    const jobIds = iterator[0];
    console.log('Total jobs are', jobIds.length);

    const batchSize = 100;
    // 把jobIds分成多个批次
    const batches = [];
    for (let i = 0; i < jobIds.length; i += batchSize) {
      batches.push(jobIds.slice(i, i + batchSize));
    }

    for (const batch of batches) {
      console.log(`Processing batch of ${batch.length} jobs`);
      // 并发请求当前批次的所有job
      const metricsPromises = batch.map(jobId => fetchJobStatus(jobId).catch(err => {
        console.error(`处理Job ${jobId}失败: ${err.message}`);
        return null; // 失败的job返回null,后续过滤掉
      }));
      const batchMetrics = (await Promise.all(metricsPromises)).filter(Boolean); // 过滤掉失败的

      if (batchMetrics.length > 0) {
        console.log(`Writing ${batchMetrics.length} records to database`);
        // await writeToInfluxDB(batchMetrics);
      }
    }

    console.log('所有Job处理完成');
  } catch (err) {
    console.error('读取文件失败:', err.message);
  }
}

这个版本会把job分成多个批次,每个批次并发请求,处理速度更快,但要注意你的服务能承受的并发量,避免请求过多导致被限流。

内容的提问来源于stack exchange,提问作者user14435933

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.29 06:12:36