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);
关键修复说明
- 进程信号监听:捕获
SIGINT(Ctrl+C)和SIGTERM系统信号,统一处理终止逻辑 - 流管控:终止时暂停Scram Jet流、销毁MySQL查询流,彻底停止数据流入和新任务启动
- 资源释放:及时释放数据库连接,避免资源泄漏
- 异步任务管控:
- 双重检查
stopSending和isShuttingDown,阻止新任务启动 - 可选使用Axios CancelToken取消正在进行的API请求,减少无效请求
- 双重检查
- 幂等终止:通过
isShuttingDown标记确保终止流程只执行一次
内容的提问来源于stack exchange,提问作者Benson OO
相关产品推荐
相关产品推荐

