Promise 并发调度器 Scheduler

实现一个异步任务调度器,保证任意时刻最多只有 maxConcurrency 个任务正在执行。一个任务结束后,立即从等待队列中启动下一个任务。

题目

实现一个 Scheduler 类:

  • 构造函数接收最大并发数;
  • add(task) 添加异步任务;
  • task 是一个返回 Promise 的函数;
  • add() 返回一个 Promise,其结果与当前任务一致;
  • 单个任务失败不能阻塞后续任务。
const scheduler = new Scheduler(2);

const timeout = (milliseconds) =>
  new Promise((resolve) => {
    setTimeout(resolve, milliseconds);
  });

const addTask = (milliseconds, value) => {
  scheduler.add(() => timeout(milliseconds)).then(() => {
    console.log(value);
  });
};

addTask(1000, '1');
addTask(500, '2');
addTask(300, '3');
addTask(400, '4');

// 最大并发数为 2,输出顺序:2 3 1 4

实现

class Scheduler {
  constructor(maxConcurrency) {
    if (!Number.isInteger(maxConcurrency) || maxConcurrency <= 0) {
      throw new RangeError('maxConcurrency must be a positive integer');
    }

    this.maxConcurrency = maxConcurrency;
    this.runningCount = 0;
    this.queue = [];
  }

  add(task) {
    if (typeof task !== 'function') {
      return Promise.reject(new TypeError('task must be a function'));
    }

    return new Promise((resolve, reject) => {
      this.queue.push({ task, resolve, reject });
      this.runNext();
    });
  }

  runNext() {
    while (
      this.runningCount < this.maxConcurrency &&
      this.queue.length > 0
    ) {
      const { task, resolve, reject } = this.queue.shift();
      this.runningCount += 1;

      Promise.resolve()
        .then(task)
        .then(resolve, reject)
        .finally(() => {
          this.runningCount -= 1;
          this.runNext();
        });
    }
  }
}

核心思路

调度器维护两个状态:

  • queue:尚未开始的任务队列;
  • runningCount:当前正在执行的任务数量。

添加任务时,先将任务放入队列,再尝试调度。只要运行数量没有达到上限,就从队首取出任务执行。任务无论成功还是失败,都在 finally 中释放一个并发名额,然后继续调度。

这里的队列只能保证任务按照加入顺序开始,不能保证按照加入顺序完成。任务完成顺序取决于各自耗时。

为什么必须传入函数

正确写法:

scheduler.add(() => fetch('/api/user'));

错误写法:

scheduler.add(fetch('/api/user'));

fetch() 调用后请求已经开始。即使随后把 Promise 放入队列,调度器也无法控制它的启动时机。将请求包装成函数,才能把任务的执行权交给调度器。

为什么使用 Promise.resolve().then(task)

不能只写:

const promise = task();

因为 task() 可能同步抛出异常,导致后面的 Promise 链和并发名额释放逻辑无法执行。

Promise.resolve().then(task)

可以统一处理三种情况:

  • 返回 Promise;
  • 返回普通值;
  • 同步抛出异常。

失败任务为什么不会阻塞队列

.then(resolve, reject)
.finally(() => {
  this.runningCount -= 1;
  this.runNext();
});

当前任务失败时,add() 返回的 Promise 会被拒绝,但 finally 仍然执行,因此并发名额会正常释放,等待队列也会继续运行。

调用方仍然需要处理失败,否则可能出现未处理的 Promise rejection:

scheduler
  .add(() => fetch('/api/user'))
  .catch((error) => console.error(error));

复杂度

上述实现使用数组的 shift() 取出队首任务,每次操作最坏为 O(n)。面试题通常可以接受;如果任务量很大,可以使用游标读取数组,或者实现真正的队列,将出队操作优化为 O(1)。

空间复杂度为 O(n),主要用于保存等待执行的任务。

高频追问

JavaScript 是单线程,为什么还存在并发?

这里的并发不是多个 JavaScript 任务同时占用主线程计算,而是多个网络请求、定时器或其他异步操作处于进行中。事件循环负责在异步操作完成后执行对应回调。

为什么不用 Promise.all?

Promise.all() 会立即订阅传入的 Promise。如果 Promise 在创建时就启动了请求,所有请求会一次性发出,无法限制最大并发数。

如何保证结果顺序?

Scheduler 的每次 add() 都返回当前任务对应的 Promise。批量处理时可以保存这些 Promise,最后使用 Promise.all(),其结果顺序与传入 Promise 的顺序一致:

const tasks = urls.map((url) => scheduler.add(() => fetch(url)));
const responses = await Promise.all(tasks);

如何取消正在执行的请求?

调度器只能决定任务何时开始,不能自动终止已经运行的任务。网络请求需要结合 AbortController,并将 signal 传给 fetch()

如何增加优先级?

为队列元素增加 priority,调度时优先取出优先级更高的任务。优先级相同时,还应根据加入顺序保持稳定排序。

面试总结

这道题的本质是“等待队列 + 运行计数器”:

  1. 添加任务时入队;
  2. 有空闲名额时启动队首任务;
  3. 任务完成或失败后释放名额;
  4. 继续启动下一个等待任务。

回答时要特别强调:调度器接收的是任务函数,而不是已经开始执行的 Promise。