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

Node.js中MySQL流式查询遇强制关闭仍有API调用的修复方法

问题修复:流式处理强制关闭后仍有API调用执行的解决方案

问题根源

当前实现的核心问题:

  • 仅靠stopSending变量只能阻止新API调用启动,但已经进入异步执行阶段的请求(已发起的API调用、数据库更新)会继续执行
  • MySQL查询流未在关闭时中断,仍可能有剩余数据写入newDataStream
  • 缺少进程强制关闭时的优雅终止逻辑,流和异步任务未被统一管控

修复方案

核心处理步骤

  • 监听进程强制关闭信号(SIGINT/SIGTERM),触发统一终止流程
  • 中断MySQL查询流,停止读取新数据
  • 暂停Scram Jet DataStream,阻止新异步任务启动
  • 等待正在执行的异步任务完成后(或按需强制终止),释放资源并退出进程

修复后的代码

const axios = require('axios');
const { DataStream } = require('scramjet');

let stopSending = false;
let isShuttingDown = false;
let newDataStream = new DataStream();
let numbersStream = null;
let dbConnection = null;

// 监听进程强制关闭信号
process.on('SIGINT', handleShutdown);
process.on('SIGTERM', handleShutdown);

function handleShutdown() {
  if (isShuttingDown) return;
  isShuttingDown = true;

  console.log('开始终止流程...');
  stopSending = true;

  // 暂停Scram Jet流,阻止新任务启动
  newDataStream.pause();

  // 中断MySQL查询流
  if (numbersStream) {
    numbersStream.destroy();
    numbersStream = null;
  }

  // 释放数据库连接
  if (dbConnection) {
    dbConnection.release();
    dbConnection = null;
  }

  // 等待所有异步任务完成后退出(若需强制退出可直接调用process.exit())
  newDataStream.on('finish', () => {
    console.log('所有任务已完成,进程退出');
    process.exit(0);
  });
}

pool.getConnection((err, connection) => {
  if (err) {
    console.error('获取数据库连接失败:', err);
    return;
  }
  dbConnection = connection;

  numbersStream = connection.query(`SELECT id,number from ${msisdnTableName} Where status='Pending'`).stream({
    highWaterMark: 5
  });

  numbersStream.on('data', (data) => {
    if (!stopSending) {
      newDataStream.write(data);
    }
  });

  numbersStream.on('end', () => {
    console.log('MySQL查询流已结束');
    connection.release();
    dbConnection = null;
    newDataStream.end();
  });

  numbersStream.on('error', (err) => {
    console.error('MySQL流错误:', err);
    connection.release();
    dbConnection = null;
    newDataStream.end();
  });
});

newDataStream.map(async (numberObj) => {
  // 双重检查,确保终止时不再启动新任务
  if (stopSending || isShuttingDown) return;

  // 可选:使用CancelToken支持取消正在进行的API请求
  const source = axios.CancelToken.source();
  try {
    const response = await axios.get("http://localhost:8081/data", {
      params: { number: numberObj.number },
      cancelToken: source.token
    });

    if (response.status === 200) {
      await mysqlUpdateStatus(numberObj.id, 'Sent');
    } else {
      await mysqlUpdateStatus(numberObj.id, 'Failed');
    }
  } catch (err) {
    if (axios.isCancel(err)) {
      console.log(`请求已取消: ${numberObj.number}`);
    } else {
      console.error(`处理数据失败: ${numberObj.id}`, err);
      await mysqlUpdateStatus(numberObj.id, 'Failed');
    }
  } finally {
    // 若正在关闭,取消未完成的请求
    if (isShuttingDown) {
      source.cancel('进程正在关闭,取消请求');
    }
  }
})
.on('end', () => {
  console.log('所有数据处理完成');
})
.on('error', (err) => {
  console.error('DataStream处理错误:', err);
});

// 模拟强制关闭(实际场景由进程信号触发)
setTimeout(() => {
  console.log('模拟触发强制关闭');
  handleShutdown();
}, 1000);

关键修复说明

  1. 进程信号监听:捕获SIGINT(Ctrl+C)和SIGTERM系统信号,统一处理终止逻辑
  2. 流管控:终止时暂停Scram Jet流、销毁MySQL查询流,彻底停止数据流入和新任务启动
  3. 资源释放:及时释放数据库连接,避免资源泄漏
  4. 异步任务管控:
    • 双重检查stopSending和isShuttingDown,阻止新任务启动
    • 可选使用Axios CancelToken取消正在进行的API请求,减少无效请求
  5. 幂等终止:通过isShuttingDown标记确保终止流程只执行一次

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.06 06:20:40