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

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);
  });
});

关键逻辑说明

  1. 新行检测:通过记录lastProcessedId,每次仅查询ID大于该值的行,避免重复处理。
  2. 请求生成:解析json列为Axios请求配置,结合ip字段自定义请求参数(可根据实际业务场景调整逻辑)。
  3. 错误处理:包含JSON解析失败、请求执行失败、数据库操作失败的捕获与日志输出,可按需添加重试逻辑。
  4. 资源管理:监听进程退出信号,确保数据库连接正常关闭。

优化建议

  • 若追求实时性,可替换轮询为MySQL Binlog监听(使用mysql-binlog-connector包),减少数据库查询压力。
  • 给request表的id字段添加主键索引,提升查询效率。
  • 高并发场景下,可引入消息队列(如Bull)异步处理请求,避免阻塞主线程。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 03:40:34