将模块变量转换为实例变量:async/await并发处理模块开发问询
解决方案:将模块状态改为实例变量
要解决这个问题,我们可以把原来的模块级全局变量封装到类中,这样每个实例都会拥有独立的队列、处理函数和并发配置,完全避免全局状态冲突。下面是具体的实现思路和完整代码:
核心思路
原来的代码把internalQueue、func、parallelism都挂在module对象上,属于全局共享状态——这意味着所有使用该模块的地方都会共用同一个队列,显然不符合多实例的需求。用类封装后,每个实例的状态都是独立的,你可以同时创建多个并发处理器,各自处理不同的任务流。
完整实现代码
class AsyncConcurrentProcessor { constructor(parallelism = 1) { // 实例变量:每个实例独立拥有这些状态 this.internalQueue = []; this.func = undefined; this.parallelism = parallelism; this.runningTasks = 0; // 新增:跟踪当前运行的任务数,精准控制并发 this.isProcessing = false; // 新增:标记是否已启动处理流程 } // 设置任务处理函数 setHandler(handlerFunc) { if (typeof handlerFunc !== 'function') { throw new Error('Handler must be an async function'); } this.func = handlerFunc; } // 添加任务到队列 addTask(taskData) { this.internalQueue.push(taskData); // 添加任务后自动启动处理(如果还没启动) if (!this.isProcessing) { this.startProcessing(); } } // 动态调整并发数 setParallelism(num) { if (typeof num !== 'number' || num < 1) { throw new Error('Parallelism must be a positive number'); } this.parallelism = num; // 调整后如果还有待处理任务,尝试启动更多任务 if (this.isProcessing) { this.processNext(); } } // 核心处理逻辑:递归启动任务,严格控制并发数 async processNext() { // 满足条件时持续启动新任务 while (this.internalQueue.length > 0 && this.func && this.runningTasks < this.parallelism) { this.runningTasks++; const taskData = this.internalQueue.shift(); try { await this.func(taskData); } catch (err) { // 这里可以根据需求添加错误处理,比如日志上报、任务重试 console.error('Task execution failed:', err); } finally { this.runningTasks--; // 当前任务完成后,自动尝试启动下一个 this.processNext(); } } // 所有任务处理完成后,重置状态 if (this.internalQueue.length === 0 && this.runningTasks === 0) { this.isProcessing = false; } } // 手动启动处理流程(可选,因为addTask会自动触发) startProcessing() { if (!this.func) { throw new Error('Please set a handler function first'); } this.isProcessing = true; this.processNext(); } } // 导出类,用户可以创建多个独立实例 module.exports = AsyncConcurrentProcessor;
使用示例
const AsyncProcessor = require('./your-module-path'); // 创建第一个实例,并发数设为2 const processor1 = new AsyncProcessor(2); processor1.setHandler(async (task) => { console.log(`Processing task ${task} in processor1`); await new Promise(resolve => setTimeout(resolve, 1000)); }); processor1.addTask('A'); processor1.addTask('B'); processor1.addTask('C'); // 创建第二个完全独立的实例,并发数设为1 const processor2 = new AsyncProcessor(1); processor2.setHandler(async (task) => { console.log(`Processing task ${task} in processor2`); await new Promise(resolve => setTimeout(resolve, 500)); }); processor2.addTask('X'); processor2.addTask('Y');
关键改进点
- 实例隔离:每个
AsyncConcurrentProcessor实例都有自己的队列、处理函数和并发配置,完全独立运行,不会互相干扰。 - 并发控制更精准:新增
runningTasks计数器,严格限制同时运行的任务数,避免超出设定的并发上限。 - API更直观:提供
setHandler、addTask、setParallelism等明确的方法,比直接操作模块变量更易维护。 - 自动触发处理:添加任务后自动启动处理流程,无需手动调用启动方法,使用更便捷。
内容的提问来源于stack exchange,提问作者Moshe Shmukler
相关产品推荐
相关产品推荐

