读取评价CSV调用API:行计数实现与断点续传方案
实现CSV数据API调用的断点续传与限流控制
针对你的需求,我会帮你在现有代码基础上添加行计数、断点续传和合规限流功能,确保脚本中断后能从上次停止的位置继续执行,同时不超过API的调用限制。
核心思路
- 用本地文件记录当前处理到的行号(断点续传的关键)
- 维护调用计数器,每达到300次就暂停5分钟(适配限流规则)
- 每次处理完一行就更新进度文件,避免意外中断丢失进度
修改后的完整代码
const fs = require('fs').promises; const path = require('path'); const csv = require('csvtojson'); const axios = require('axios'); const moment = require('moment'); // 配置项 const PROGRESS_FILE = path.join(__dirname, 'progress.json'); const CSV_FILE = path.join(__dirname, "/csv/review.csv"); const API_LIMIT = 300; // 每5分钟最多调用次数 const LIMIT_WAIT_TIME = 5 * 60 * 1000; // 5分钟(毫秒) const CALL_DELAY = 5000; // 每次调用后的基础延迟 async function saveProgress(lineNumber) { await fs.writeFile(PROGRESS_FILE, JSON.stringify({ lastLine: lineNumber })); } async function loadProgress() { try { const data = await fs.readFile(PROGRESS_FILE, 'utf8'); const progress = JSON.parse(data); return progress.lastLine || 0; } catch (err) { // 如果进度文件不存在,默认从第0行开始 return 0; } } async function createReview2() { try { // 加载上次处理的进度 let lastProcessedLine = await loadProgress(); console.log(`从第 ${lastProcessedLine + 1} 行开始处理...`); // 读取CSV文件 const jsonArray = await csv({ delimiter: ';', noheader: false, }).fromFile(CSV_FILE); let callCount = 0; // 从上次停止的位置开始循环 for (let i = lastProcessedLine; i < jsonArray.length; i++) { const d = jsonArray[i]; const currentLine = i + 1; // 行号从1开始计数 console.log(`正在处理第 ${currentLine} 行...`); // 构建请求体 let rev = { "review": { "createdAt": moment(d.date_added, "DD/MM/YYYY HH:mm:ss").format(), "isVisible": true, "score": d.rating, "name": d.author, "content": d.text, "language": "en", "customer": "122862444", "product": d['xx ID'] } }; try { // 执行API调用 await axios({ method: 'post', data: rev, url: 'https://xx:xx@api.xx.com/en/reviews.json', }); callCount++; console.log(`第 ${currentLine} 行处理成功,累计调用次数:${callCount}`); // 保存当前进度 await saveProgress(i); // 处理限流:每达到300次调用,暂停5分钟 if (callCount % API_LIMIT === 0) { console.log(`已达到每5分钟300次调用限制,将暂停5分钟...`); await new Promise(resolve => setTimeout(resolve, LIMIT_WAIT_TIME)); } else { // 基础延迟,避免短时间内密集调用 await new Promise(resolve => setTimeout(resolve, CALL_DELAY)); } } catch (apiErr) { console.error(`第 ${currentLine} 行处理失败:`, apiErr.message); // 可以选择重试,或者记录失败行到文件,这里暂时跳过继续下一行 // 如果需要重试,可添加重试逻辑 await saveProgress(i); // 即使失败也保存进度,避免重复处理 await new Promise(resolve => setTimeout(resolve, CALL_DELAY)); } } console.log("所有行处理完成!"); // 处理完成后可以删除进度文件,或者保留以便后续查看 await fs.unlink(PROGRESS_FILE).catch(() => {}); } catch (err) { console.error("脚本执行出错:", err.message); } } // 注意:原代码中的sleep函数需要替换为Promise版,这里已经用setTimeout封装了 createReview2();
关键功能说明
- 断点续传:通过
saveProgress和loadProgress函数读写progress.json文件,记录每次处理完的行索引,重启脚本时直接从该位置继续。 - 限流控制:维护
callCount计数器,每累计300次调用就触发5分钟的等待,严格遵守API的限流规则;同时保留了原有的5秒基础延迟,降低触发限流的风险。 - 异常处理:API调用失败时会记录错误并保存当前进度,避免重复处理失败行(如果需要重试,可以在catch块中添加重试逻辑)。
- 日志输出:添加了行号和调用次数的日志,方便跟踪处理进度。
使用注意事项
- 确保
progress.json文件所在目录有读写权限。 - 如果需要修改限流规则,直接调整
API_LIMIT和LIMIT_WAIT_TIME配置项即可。 - 如果CSV文件内容有修改(比如新增行),脚本会自动处理新增的行,不会重复处理已完成的行。
内容的提问来源于stack exchange,提问作者I. T.
相关产品推荐
相关产品推荐

