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

将模块变量转换为实例变量: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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.26 09:46:38