Node.js进程结束前如何执行异步操作并等待完成?
解决SIGINT信号下APM批量数据未持久化的问题
我正在构建一个简单的APM工具,将事件批量存储到数据库中,同时希望在进程结束、被杀死或出错时,存储批次中剩余的所有数据。当前SIGINT处理器能正常触发,能看到APM - Graceful shutdown的日志,但数据库中没有存储任何数据,进程似乎在await insertExecutionTimes(batched);完成前就已终止。
现有代码
启动APM的核心逻辑
export function startApm(): void { process.on('SIGINT', async () => { try { console.log('APM - Graceful shutdown'); console.log(`Will insert ${batched.length} entries`); await insertExecutionTimes(batched); console.log('Execution times inserted successfully'); } catch (error) { console.error('Error inserting apm entries:', error); } finally { process.exit(0); } }); setInterval(async () => { if (batched.length >= batchSize) { console.log(`Storing batch of ${batchSize}, total length ${batched.length}`); await insertExecutionTimes(batched.splice(0, batchSize)); console.log('Length after', batched.length); } }, 250); }
数据库插入逻辑
async function insertExecutionTimes(data: ApmEntry[]): Promise<void> { const trx = await db.transaction(); try { await trx.batchInsert(Tables.Apm, data); await trx.commit(); console.log(`Inserted ${data.length} records successfully`); } catch (error) { await trx.rollback(); console.error('Error inserting execution times', error); } }
解决思路
- 不要在异步信号处理器中直接强制退出进程:Node.js的信号回调不会自动等待异步操作完成,你在
finally里调用process.exit(0)会直接终止进程,跳过await insertExecutionTimes的执行。必须等所有异步操作完成后再调用process.exit。 - 标记关闭状态,避免冲突:新增一个状态变量,阻止定时器在进程关闭过程中继续处理批次,防止和关闭逻辑抢数据。
- 覆盖更多终止信号:除了SIGINT,还需要处理SIGTERM(系统kill命令触发的信号),确保不同终止场景下都能执行数据持久化。
- 确保事务逻辑完整:在插入函数中抛出错误,让上层处理器捕获,同时避免空数据的无效事务操作。
修改后的示例代码
let isShuttingDown = false; let batched: ApmEntry[] = []; const batchSize = 100; export function startApm(): void { // 处理Ctrl+C触发的SIGINT process.on('SIGINT', handleShutdown); // 处理系统kill命令触发的SIGTERM process.on('SIGTERM', handleShutdown); setInterval(async () => { // 关闭状态下停止批量处理 if (isShuttingDown || batched.length < batchSize) return; console.log(`Storing batch of ${batchSize}, total length ${batched.length}`); await insertExecutionTimes(batched.splice(0, batchSize)); console.log('Length after', batched.length); }, 250); } async function handleShutdown(): Promise<void> { // 避免重复触发关闭逻辑 if (isShuttingDown) return; isShuttingDown = true; console.log('APM - Graceful shutdown'); console.log(`Will insert ${batched.length} entries`); try { await insertExecutionTimes(batched); console.log('Execution times inserted successfully'); } catch (error) { console.error('Error inserting apm entries:', error); } finally { // 所有异步操作完成后再退出进程 process.exit(0); } } async function insertExecutionTimes(data: ApmEntry[]): Promise<void> { // 空数据直接返回,避免无效事务 if (data.length === 0) return; const trx = await db.transaction(); try { await trx.batchInsert(Tables.Apm, data); await trx.commit(); console.log(`Inserted ${data.length} records successfully`); } catch (error) { await trx.rollback(); console.error('Error inserting execution times', error); // 向上抛出错误,让上层关闭处理器捕获 throw error; } }
内容的提问来源于stack exchange,提问作者Kristi Jorgji
相关产品推荐
相关产品推荐

