基于Express JS构建插件间事件系统:同步异步事件处理及超时管控
多插件事件系统解决方案
一、同步事件的超时终止
JS主线程的同步任务无法被强制中断,唯一可行的方式是将同步任务隔离到独立执行线程中,通过终止线程停止超时任务:
Node.js 环境
使用worker_threads模块创建子线程执行同步任务,设置超时时间后终止子线程:
const { Worker, parentPort } = require('worker_threads'); function runSyncJobWithTimeout(jobInfo, timeoutMs) { return new Promise((resolve, reject) => { // 子线程执行同步任务,通过文件路径引用插件任务,避免环境丢失 const worker = new Worker(jobInfo.filePath, { workerData: { args: jobInfo.args } }); const timeout = setTimeout(() => { worker.terminate(); reject(new Error('同步任务超时')); }, timeoutMs); worker.on('message', (result) => { clearTimeout(timeout); resolve(result); worker.terminate(); }); worker.on('error', (err) => { clearTimeout(timeout); reject(err); }); }); } // 在dispatchEvent中调用 jobs.forEach(async (job) => { if (!job.isAsync) { try { await runSyncJobWithTimeout(job, 5000); // 5秒超时阈值 } catch (err) { console.error('同步任务执行失败:', err); } } });
要求插件将同步任务封装为独立文件,通过文件路径引用,而非直接传递函数,避免变量依赖丢失。
浏览器环境
使用Web Worker实现隔离执行,主线程超时后调用worker.terminate()终止任务,逻辑与Node.js类似,只需将任务代码放在Worker脚本文件中。
二、异步事件的持久化与执行
直接存储函数代码会导致依赖缺失、环境丢失,正确做法是存储任务元数据而非代码本身:
1. 数据库存储任务元数据
每个异步任务存储以下核心信息:
eventId: 关联的事件IDpluginId: 所属插件的唯一标识taskName: 插件中异步方法的名称args: 任务执行参数(需支持JSON序列化)status: 任务状态(待执行/执行中/完成/失败)timeoutMs: 任务超时阈值
示例数据库记录:
{ "eventId": "user_register", "pluginId": "sms_notifier", "taskName": "sendVerificationCode", "args": ["138xxxxxxx"], "status": "pending", "timeoutMs": 30000 }
2. CRON任务执行异步任务
CRON定时拉取待执行任务,通过全局插件管理器调用对应插件的方法:
// 全局插件注册管理器 const pluginManager = { plugins: new Map(), register(pluginId, pluginInstance) { this.plugins.set(pluginId, pluginInstance); } }; // CRON执行逻辑 async function executePendingAsyncJobs() { const pendingJobs = await fetchPendingJobsFromDB(); // 从数据库拉取待执行任务 for (const job of pendingJobs) { const plugin = pluginManager.plugins.get(job.pluginId); if (!plugin) { await updateJobStatus(job.id, 'failed', '插件未找到'); continue; } try { const task = plugin[job.taskName]; if (typeof task !== 'function') { throw new Error('插件中不存在指定任务'); } // 超时控制:用Promise.race同时等待任务完成和超时 const timeoutPromise = new Promise((_, reject) => setTimeout(() => reject(new Error('异步任务超时')), job.timeoutMs) ); await Promise.race([task(...job.args), timeoutPromise]); await updateJobStatus(job.id, 'completed'); } catch (err) { await updateJobStatus(job.id, 'failed', err.message); } } }
3. 插件开发规范
要求插件必须:
- 暴露可调用的异步方法(返回Promise)
- 启动时通过
pluginManager完成注册 - 方法参数仅使用JSON可序列化类型(避免传递复杂对象)
三、改造原dispatchEvent函数
基于上述方案优化原事件分发逻辑:
async function dispatchEvent(eventId) { const jobs = getJobsForEvent(eventId); for (const job of jobs) { if (job.isAsync) { // 存储任务元数据到数据库 await saveAsyncJobToDB({ eventId, pluginId: job.pluginId, taskName: job.taskName, args: job.args, timeoutMs: job.timeoutMs || 30000 }); } else { // 执行带超时的同步任务 try { await runSyncJobWithTimeout(job, job.timeoutMs || 5000); } catch (err) { console.error(`事件${eventId}的同步任务执行失败:`, err); } } } }
内容的提问来源于stack exchange,提问作者Faheem Anis
相关产品推荐
相关产品推荐

