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

Node.js API实现CSV上传、多接口处理及更新文件返回方案问询

问题

我已经编写了一个通过调用多个API获取信息来更新CSV文件的函数,现在希望开发一个Node.js API,支持上传CSV文件并调用该函数处理,最终将更新后的文件返回给客户端。请问该如何实现?还有其他可行方案吗?

以下是我的核心处理函数:

async function processCSV(filePath) {
  try {
    const rows = []
    const uniqueReferenceNumbers = new Set()

    fs.createReadStream(filePath)
      .pipe(csv())
      .on('data', (row) => {
        const referenceNumber = row['Reference Number']
        console.log('读取CSV,参考编号:', referenceNumber)
        // 处理前检查参考编号是否唯一
        if (!uniqueReferenceNumbers.has(referenceNumber)) {
          uniqueReferenceNumbers.add(referenceNumber)
          rows.push(row)
        }
      })
      .on('end', async () => {
        if (rows.length === 0) {
          console.log('CSV文件中未找到唯一参考编号。')
          return
        }

        // 假设CSV中所有参考编号的品牌名称一致
        const brandName = rows[0]['Brand']

        // 第一次API调用获取brand_uuid
        const brandApiResponse = await makeAPICallWithRetry(
          `https://api/v2/search/brand?q=${brandName}`
        )
        const brandUuid = brandApiResponse.data.data[0].uuid
        console.log('品牌UUID:', brandUuid)

        // 为每个唯一参考编号进行第二、第三、第四、第五次API调用
        const updatedRows = [] // 在内存中存储修改后的行
        for (const row of rows) {
          const referenceNumber = row['Reference Number']

          try {
            // 使用获取到的brand_uuid和reference_number进行第二次API调用
            const secondApiResponse = await makeAPICallWithRetry(
              `https://api/v2/search/watch?q=${referenceNumber}&brand_uuid=${brandUuid}`
            )

            // 从第二次API响应中提取watch_uuid
            const watchUuid = secondApiResponse.data.data[0].uuid
            console.log('腕表UUID:', watchUuid)
            // 第三次API调用获取产品信息
            const thirdApiResponse = await makeAPICallWithRetry(
              `https://api/v2/watch/info?uuid=${watchUuid}`
            )

            // 第四次API调用获取产品规格
            const fourthApiResponse = await makeAPICallWithRetry(
              `https://api/v2/watch/specs?uuid=${watchUuid}`
            )
            // 第五次API调用获取价格历史
            const fifthApiResponse = await makeAPICallWithRetry(
              `https://api/v2/watch/price?uuid=${watchUuid}`
            )
            row['产品信息'] = thirdApiResponse
              ? JSON.stringify(thirdApiResponse.data, null, 2)
              : ''
            row['产品规格'] = fourthApiResponse
              ? JSON.stringify(fourthApiResponse.data, null, 2)
              : ''
            row['价格历史'] = fifthApiResponse
              ? JSON.stringify(fifthApiResponse.data, null, 2)
              : console.log(
                  `正在更新参考编号${referenceNumber},信息:${row['产品信息']},规格:${row['产品规格']},价格历史:${row['价格历史']}`
                )
            // 将修改后的行加入结果集
            updatedRows.push(row)
          } catch (error) {
            // 适当处理错误,比如日志记录、跳过该行等
            console.error(
              `处理参考编号${referenceNumber}时出错:${error.message}`
            )
            // 如果超过最大重试次数,终止整个流程
            if (
              error.message.includes('超过最大重试次数')
            ) {
              // 在抛出错误前将已累积的数据写入CSV
              if (updatedRows.length > 0) {
                const csvData = Papa.unparse(updatedRows, { header: true })
                fs.writeFileSync(filePath, csvData)
              }
              throw error
            }
          }
        }

        // 将修改后的行更新到CSV文件
        const csvData = Papa.unparse(rows, { header: true })
        fs.writeFileSync(filePath, csvData)

        // 处理完成,可以对结果进行后续操作
        console.log('API调用完成,CSV文件已更新')
      })
  } catch (error) {
    // 适当处理错误
    console.error('处理CSV时出错:', error.message)
  }
}
实现方案与替代方案

一、基础Node.js API实现步骤

1. 安装依赖

首先安装必要的包:

npm install express multer papaparse
  • express:搭建HTTP服务
  • multer:处理文件上传
  • papaparse:解析和生成CSV

2. 改造processCSV函数

原函数依赖本地文件路径,改成直接处理CSV内容字符串(避免本地文件IO,提升效率):

const Papa = require('papaparse');

async function processCSVContent(csvContent) {
  return new Promise((resolve, reject) => {
    const rows = [];
    const uniqueReferenceNumbers = new Set();

    // 解析CSV内容
    Papa.parse(csvContent, {
      header: true,
      step: (result) => {
        const row = result.data;
        const referenceNumber = row['Reference Number'];
        console.log('读取CSV,参考编号:', referenceNumber);
        if (!uniqueReferenceNumbers.has(referenceNumber)) {
          uniqueReferenceNumbers.add(referenceNumber);
          rows.push(row);
        }
      },
      complete: async () => {
        if (rows.length === 0) {
          reject(new Error('CSV文件中未找到唯一参考编号'));
          return;
        }

        try {
          const brandName = rows[0]['Brand'];
          const brandApiResponse = await makeAPICallWithRetry(
            `https://api/v2/search/brand?q=${brandName}`
          );
          const brandUuid = brandApiResponse.data.data[0].uuid;
          console.log('品牌UUID:', brandUuid);

          const updatedRows = [];
          for (const row of rows) {
            const referenceNumber = row['Reference Number'];
            try {
              const secondApiResponse = await makeAPICallWithRetry(
                `https://api/v2/search/watch?q=${referenceNumber}&brand_uuid=${brandUuid}`
              );
              const watchUuid = secondApiResponse.data.data[0].uuid;

              // 并行调用多个API,提升处理速度
              const [thirdApiResponse, fourthApiResponse, fifthApiResponse] = await Promise.all([
                makeAPICallWithRetry(`https://api/v2/watch/info?uuid=${watchUuid}`),
                makeAPICallWithRetry(`https://api/v2/watch/specs?uuid=${watchUuid}`),
                makeAPICallWithRetry(`https://api/v2/watch/price?uuid=${watchUuid}`)
              ]);

              row['Product Info'] = thirdApiResponse ? JSON.stringify(thirdApiResponse.data, null, 2) : '';
              row['Product Specs'] = fourthApiResponse ? JSON.stringify(fourthApiResponse.data, null, 2) : '';
              row['Price History'] = fifthApiResponse ? JSON.stringify(fifthApiResponse.data, null, 2) : '';
              updatedRows.push(row);
            } catch (error) {
              console.error(`处理参考编号${referenceNumber}出错:${error.message}`);
              if (error.message.includes('Exceeded the maximum number of retries')) {
                // 返回已处理的结果,终止流程
                const csvData = Papa.unparse(updatedRows, { header: true });
                resolve(csvData);
                return;
              }
              // 跳过出错行,继续处理其他行
              continue;
            }
          }

          const csvData = Papa.unparse(updatedRows, { header: true });
          resolve(csvData);
        } catch (error) {
          reject(error);
        }
      },
      error: (error) => {
        reject(new Error(`解析CSV失败:${error.message}`));
      }
    });
  });
}

优化点说明:

  • 用Papa.parse直接处理字符串,避免本地文件IO
  • 用Promise.all并行调用多个API,缩短处理时间
  • 用Promise封装流程,适配Express异步接口
  • 修正原函数最终写入rows而非updatedRows的错误

3. 搭建Express API服务

const express = require('express');
const multer = require('multer');
const app = express();
const port = 3000;

// 配置multer,将上传文件存储在内存中
const storage = multer.memoryStorage();
const upload = multer({ storage: storage });

// 上传处理接口
app.post('/process-csv', upload.single('csvFile'), async (req, res) => {
  try {
    if (!req.file) {
      return res.status(400).send('请上传CSV文件');
    }

    // 将上传的Buffer转为字符串
    const csvContent = req.file.buffer.toString('utf8');
    const processedCsv = await processCSVContent(csvContent);

    // 设置响应头,返回CSV文件
    res.setHeader('Content-Type', 'text/csv');
    res.setHeader('Content-Disposition', `attachment; filename=processed_${Date.now()}.csv`);
    res.send(processedCsv);
  } catch (error) {
    console.error('接口处理出错:', error);
    res.status(500).send(`处理失败:${error.message}`);
  }
});

app.listen(port, () => {
  console.log(`服务运行在 http://localhost:${port}`);
});

二、其他可行方案

1. 流式处理大CSV文件

如果CSV文件过大(几十MB以上),内存处理会占用过多资源,可采用流式解析+流式生成:

  • 用Papa.parse的流式API逐行读取
  • 处理完一行就写入响应流,无需将所有行存内存
  • 适合超大文件场景,降低内存占用

2. 异步队列处理

如果API调用耗时较长(比如处理大量行),同步接口易超时,可引入消息队列:

  • 用bull或bee-queue搭建任务队列
  • 上传文件后,将任务加入队列,立即返回任务ID给客户端
  • 后台进程处理队列任务,完成后存储结果
  • 客户端通过任务ID轮询获取处理状态和结果

3. Serverless部署

如果无需长期维护服务器,可采用Serverless平台:

  • 比如Vercel、AWS Lambda、阿里云函数计算
  • 上传接口作为Serverless函数,自动扩缩容
  • 适合低频次、突发型请求场景,节省服务器成本

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 17:44:51