NodeJS后台脚本开发:监听MySQL新行并发起Axios请求后删行
NodeJS脚本实现MySQL新行检测、Axios请求与行删除
依赖安装
首先安装所需的npm包:
npm install mysql async axios
完整实现代码
const async = require("async"); const mysql = require('mysql'); const axios = require('axios'); // 数据库连接配置 const connection = mysql.createConnection({ host: 'localhost', user: 'root', password: '', database: 'x' }); // 记录最后处理的请求ID,避免重复处理 let lastProcessedId = 0; // 处理单条请求行 function processRequestRow(row, callback) { try { // 解析json列的请求配置(需确保json格式合法) const requestConfig = JSON.parse(row.json); // 根据ip字段配置请求(示例:设置请求的目标地址或自定义头) if (row.ip) { // 场景1:ip作为目标服务地址,设置baseURL requestConfig.baseURL = `http://${row.ip}`; // 场景2:将ip作为请求来源头,可替换为实际需求 // requestConfig.headers = { ...requestConfig.headers, 'X-Forwarded-For': row.ip }; } // 发送Axios请求 axios(requestConfig) .then(response => { console.log(`请求ID ${row.id} 执行成功,响应状态码:${response.status}`); // 请求成功后删除该行 connection.query('DELETE FROM request WHERE id = ?', [row.id], (err) => { if (err) { console.error(`删除ID ${row.id} 失败:`, err); return callback(err); } console.log(`已清理ID ${row.id} 的请求记录`); callback(null); }); }) .catch(error => { console.error(`请求ID ${row.id} 执行失败:`, error.message); // 可选逻辑:如果需要重试失败请求,不要删除该行,留到下次轮询处理 callback(error); }); } catch (parseErr) { console.error(`解析ID ${row.id} 的JSON配置失败:`, parseErr.message); callback(parseErr); } } // 轮询检测新请求行 function checkNewRequests() { // 查询上次处理之后的新行 connection.query('SELECT * FROM request WHERE id > ? ORDER BY id ASC', [lastProcessedId], (err, rows) => { if (err) { console.error('查询新请求失败:', err); return; } if (rows.length === 0) { console.log('暂无新请求行'); return; } // 串行处理请求,避免并发量过高 async.eachSeries(rows, (row, cb) => { processRequestRow(row, (processErr) => { // 更新最后处理的ID(若需重试失败请求,可仅在成功时更新) lastProcessedId = row.id; cb(processErr); }); }, (seriesErr) => { if (seriesErr) { console.error('批量处理请求出错:', seriesErr); } else { console.log('本次新请求处理完成'); } }); }); } // 启动服务 connection.connect((err) => { if (err) { console.error('数据库连接失败:', err); process.exit(1); } console.log('数据库连接成功'); // 初始执行一次,之后每隔5秒轮询(可根据业务调整间隔) checkNewRequests(); setInterval(checkNewRequests, 5000); }); // 捕获退出信号,关闭数据库连接 process.on('SIGINT', () => { connection.end((err) => { if (err) console.error('关闭数据库连接失败:', err); console.log('脚本退出,数据库连接已关闭'); process.exit(0); }); });
关键逻辑说明
- 新行检测:通过记录
lastProcessedId,每次仅查询ID大于该值的行,避免重复处理。 - 请求生成:解析
json列为Axios请求配置,结合ip字段自定义请求参数(可根据实际业务场景调整逻辑)。 - 错误处理:包含JSON解析失败、请求执行失败、数据库操作失败的捕获与日志输出,可按需添加重试逻辑。
- 资源管理:监听进程退出信号,确保数据库连接正常关闭。
优化建议
- 若追求实时性,可替换轮询为MySQL Binlog监听(使用
mysql-binlog-connector包),减少数据库查询压力。 - 给
request表的id字段添加主键索引,提升查询效率。 - 高并发场景下,可引入消息队列(如Bull)异步处理请求,避免阻塞主线程。
内容的提问来源于stack exchange,提问作者rapstar2004
相关产品推荐
相关产品推荐

