先说说我为什么写这个玩意儿。上个月接了个私活,要给一个旧系统写批量数据导入。前端需要一次性提交几百个请求,但是后端接口只允许同时处理5个并发,多了就报错。更坑的是,偶尔网络抖动,请求失败就得重试。我一开始用了Promise.all,直接被后端骂死。后来用了个简单的池子,但没法取消,用户点个“停止”按钮一点反应都没有。你说这像话吗?
所以我腾出一下午,手写了一个异步调度器。今天你就跟着我一起把这个东西从零搭起来。别怕,代码不超过100行,但能解决80%的异步控制需求。
先列需求,别急着写代码
我脑子里混乱的时候是写不出好代码的,所以先把要什么写下来:
- 能控制最大并发数,比如同一时刻最多跑5个任务。
- 任务是一个返回Promise的函数,不是已经发起的Promise。
- 支持取消:用户点了取消,正在执行的任务如果能中断就拉倒,还在排队的任务统统不让跑了。
- 支持自动重试:单个任务失败后,自动重新执行,最多重试N次。
- 出错的终极结果要返回reject,不能让异常被吞掉。
光这几个,用现成的库比如p-limit能搞定前两个,但是后面俩就得自己融合。我干脆全手写,把控每一个细节。
第一步:设计调用方式
大概长这样:
const scheduler = new Scheduler({
concurrency: 5,
retry: 2
});
// 提交一个任务
const promise = scheduler.add(() => fetch('/api/data'), {
signal: someAbortSignal
});
// 等结果
promise.then(console.log).catch(console.error);
这里的retry: 2指的是任务失败后额外重试2次,一共最多执行3次。signal可以随时取消,取消后这个任务即使已经在执行,如果它内部接受signal,也能停止。
我还想支持优先级,但思来想去,这个业务场景暂时用不上,加了反而变复杂。咱们先把基础打牢。
第二步:搭出队列和运行池
调度器内部无非就是一个数组保存等待的任务,再统计正在运行的任务数量。当运行数量小于并发上限时,就从队列里拿一个任务出来跑。
class Scheduler {
constructor(options = {}) {
this.concurrency = options.concurrency || 5;
this.retry = options.retry || 0;
this.retryDelay = options.retryDelay || 0;
this.queue = [];
this.runningCount = 0;
}
add(taskFn, { signal } = {}) {
return new Promise((resolve, reject) => {
const job = {
taskFn,
signal,
resolve,
reject,
retriesLeft: this.retry
};
this.queue.push(job);
// 如果有signal传入,该信号取消时,直接把job标记为已取消
if (signal) {
signal.addEventListener('abort', () => {
this._disposeJob(job, new Error('task aborted'));
}, { once: true });
}
this._tick();
});
}
_disposeJob(job, err) {
const idx = this.queue.indexOf(job);
if (idx !== -1) {
this.queue.splice(idx, 1);
}
job.reject(err);
}
_tick() {
while (this.runningCount < this.concurrency && this.queue.length > 0) {
const job = this.queue.shift();
this._executeJob(job);
}
}
async _executeJob(job) {
this.runningCount++;
try {
const result = await this._runWithRetry(job);
job.resolve(result);
} catch (err) {
job.reject(err);
} finally {
this.runningCount--;
this._tick();
}
}
}
这里有个细节:为什么在add()里返回了一个新的Promise而不是直接用taskFn的Promise?因为我们要把resolve和reject挂在自定义job上,让重试和取消都在这个维度操作。
_tick()的意思是“没满就继续跑”。每当一个任务完成,运行数减一,再触发_tick(),队列里下一个任务就能补上来。
第三步:实现重试
重试逻辑不能简单包一层for循环。因为失败后可能需要延迟,而且如果任务被取消了,不应该再重试。我这么写:
async _runWithRetry(job) {
while (true) {
try {
// 如果任务在开始前就被取消,直接抛错
if (job.signal?.aborted) {
throw new Error('task aborted');
}
// 把AbortSignal传给任务函数,让网络请求等支持取消的内部逻辑能中断
const result = await job.taskFn(job.signal);
return result;
} catch (err) {
// 如果已经被取消,直接抛出,不再重试
if (err?.name === 'AbortError' || job.signal?.aborted) {
throw err;
}
if (job.retriesLeft > 0) {
job.retriesLeft--;
if (this.retryDelay > 0) {
await new Promise(resolve => setTimeout(resolve, this.retryDelay));
}
// continue 触发下一轮循环,重新执行任务
} else {
throw err;
}
}
}
}
就是这么简单。用while(true)包住任务,每失败一次就扣一次重试次数,扣完就throw。同时注意检查signal.aborted,如果用户取消了,就不再浪费时间重试。
为什么要把signal传给taskFn?因为fetch支持signal参数。如果你传入的signal是一个AbortController的信号,那么fetch才会真正中断。不然你只是“单方面宣布取消”,底层的请求还在跑,资源就浪费了。
第四步:完善取消逻辑
用户通过AbortSignal来取消任务,但是同一时刻会有多个任务共用一个信号吗?有可能。如果用户在调度器外创建了一个AbortController,然后把这个controller.signal传给多个任务。当controller.abort()被调用,所有任务都应该取消。
我的写法是:job.signal监听abort事件,然后调用_disposeJob。但是这里有个隐患:_disposeJob会直接reject掉job,然而该job可能已经不在队列里了(它正在执行中)。所以_disposeJob里必须同时考虑两种情况:正在排队的任务,或正在执行的任务。
针对正在执行的任务,直接reject会出问题吗?不会,因为我们已经在_executeJob里把job.resolve和job.reject传出去了。如果job被reject了,等_executeJob里try/catch又去resolve,就会变成“resolve一个已经settled的Promise”,无效。但要注意不能重复执行taskFn。所以我们还要在job上加一个标记:_disposed。如果任务被dispose,那么_executeJob应该检测到这个标记,不再启动重试循环。
我调整一下代码,加个状态字段:
add(taskFn, { signal } = {}) {
return new Promise((resolve, reject) => {
const job = {
taskFn,
signal,
resolve,
reject,
retriesLeft: this.retry,
disposed: false
};
this.queue.push(job);
if (signal) {
signal.addEventListener('abort', () => {
this._disposeJob(job, new Error('task aborted'));
}, { once: true });
}
this._tick();
});
}
_disposeJob(job, err) {
if (job.disposed) return;
job.disposed = true;
const idx = this.queue.indexOf(job);
if (idx !== -1) {
this.queue.splice(idx, 1);
}
job.reject(err);
}
然后_executeJob里的_runWithRetry要额外检查job.disposed,一旦发现这个标记,就不要再执行了。
async _runWithRetry(job) {
while (!job.disposed) {
try {
if (job.signal?.aborted) throw new Error('task aborted');
return await job.taskFn(job.signal);
} catch (err) {
if (err?.name === 'AbortError' || job.signal?.aborted || job.disposed) throw err;
if (job.retriesLeft > 0) {
job.retriesLeft--;
if (this.retryDelay > 0) await new Promise(r => setTimeout(r, this.retryDelay));
} else {
throw err;
}
}
}
throw new Error('task aborted');
}
你可能会问:如果任务已经执行完,但finally块还没跑,这时abort事件来了,会怎样?因为Promise已经settle,job.reject不会产生效果,但job.disposed会被置为true。这没有影响,顶多让_executeJob里的finally更快不用管。
第五步:整合代码,跑一个实例
下面我把完整代码贴出来,加上注释。
class Scheduler {
constructor({ concurrency = 5, retry = 0, retryDelay = 0 } = {}) {
this.concurrency = concurrency;
this.retry = retry;
this.retryDelay = retryDelay;
this.queue = [];
this.runningCount = 0;
}
add(taskFn, { signal } = {}) {
return new Promise((resolve, reject) => {
const job = {
taskFn,
signal,
resolve,
reject,
retriesLeft: this.retry,
disposed: false
};
this.queue.push(job);
if (signal) {
signal.addEventListener('abort', () => {
this._disposeJob(job, new Error('task aborted'));
}, { once: true });
}
this._tick();
});
}
_disposeJob(job, err) {
if (job.disposed) return;
job.disposed = true;
const idx = this.queue.indexOf(job);
if (idx !== -1) {
this.queue.splice(idx, 1);
}
job.reject(err);
}
_tick() {
while (this.runningCount < this.concurrency && this.queue.length > 0) {
const job = this.queue.shift();
this._executeJob(job);
}
}
async _executeJob(job) {
this.runningCount++;
try {
const result = await this._runWithRetry(job);
job.resolve(result);
} catch (err) {
job.reject(err);
} finally {
this.runningCount--;
this._tick();
}
}
async _runWithRetry(job) {
while (!job.disposed) {
try {
if (job.signal?.aborted) throw new Error('task aborted');
return await job.taskFn(job.signal);
} catch (err) {
if (err?.name === 'AbortError' || job.signal?.aborted || job.disposed) throw err;
if (job.retriesLeft > 0) {
job.retriesLeft--;
if (this.retryDelay > 0) await new Promise(resolve => setTimeout(resolve, this.retryDelay));
} else {
throw err;
}
}
}
throw new Error('task aborted');
}
}
实战:批量请求图片,支持取消和重试
我用这个调度器模拟一个批量上传图片缩略图的场景。假设有10个图片URL需要获取大小,每次并发3个,失败重试1次。如果用户点击了“取消”按钮,所有未完成的任务都终止。
function fetchImageSize(url, signal) {
// 模拟异步请求
return new Promise((resolve, reject) => {
if (signal?.aborted) {
reject(new Error('aborted'));
return;
}
const timer = setTimeout(() => {
if (Math.random() < 0.3) {
reject(new Error('network error'));
} else {
resolve({ url, size: Math.floor(Math.random() * 10000) });
}
}, 500 + Math.random() * 500);
signal?.addEventListener('abort', () => {
clearTimeout(timer);
reject(new Error('aborted'));
});
});
}
const scheduler = new Scheduler({
concurrency: 3,
retry: 1,
retryDelay: 200
});
const urls = Array.from({ length: 10 }, (_, i) => `https://picsum.photos/seed/${i}/200/300`);
const controller = new AbortController();
const tasks = urls.map(url => scheduler.add(() => fetchImageSize(url, controller.signal), { signal: controller.signal }));
// 3秒后自动取消,模拟用户点击“停止”
setTimeout(() => controller.abort(), 3000);
Promise.allSettled(tasks).then(results => {
const fulfilled = results.filter(r => r.status === 'fulfilled').length;
const rejected = results.filter(r => r.status === 'rejected').length;
console.log(`成功: ${fulfilled}, 失败/取消: ${rejected}`);
});
跑一下,你会发现已经开始执行的任务如果还没结束,会被abort中断;排队的任务则直接取消。而且因为重试机制,有些网络错误会触发第二次请求。这个调度器已经能覆盖大多数业务场景了。
再聊点细节:为什么不用“闲时补充”那种递归写法
很多人写并发池喜欢用for循环+递归,每完成一个就取下一个。我这种队列+_tick的方式与之相比,好处是不用维护一个游标,取消时候从数组里删掉很方便。但坏处是每个任务完成都会调用_tick,如果任务数量多且每个任务很快,可能会产生微小的性能开销。实际上这种开销比一次网络请求小几个数量级,完全无所谓。
还有一点:add()里的taskFn应该是返回Promise的函数,不能是async函数直接调用。因为async函数返回Promise,其实也符合要求。如果你想传一个已经调用过的Promise,那就不行,因为Promise已经开始了,调度器无法控制它。所以记住:一定传“函数”,不是“Promise对象”。
写到这里,做个总结吧
这个调度器是我临时手撸的,没有依赖任何库,也就几十行。你如果理解了这个,以后碰到“控制并发”“失败重试”“取消任务”的需求,可以直接复制代码改成自己的。如果还想加优先级,那就把队列改成二叉堆,或者用更简单的按priority排序。但核心调度逻辑不会变。
最后,JavaScript在处理异步时的灵活度真的很高,有时候没必要一上来就引入RxJS那种大而全的工具。自己动手写个定制化的小工具,反而能让业务代码更干净,也让你对事件循环和Promise的理解更深刻。下次面试官问“怎么控制并发数”,你可以不只是背出p-limit的用法了,而是直接画出一个调度器来。

