JavaScript并发控制实战:手写带优先级与超时的任务调度器

上周做一个数据采集后台,前端要一次性往某个接口提交几百条数据。浏览器最多同时开 6 个 TCP 连接,接口那边也扛不住瞬间的并发压力,直接 429。后来我写了一个带并发上限的任务调度器,反而把整个流程跑顺了。

相信不少人遇到过类似场景:批量上传图片、批量拉取详情、前端分片处理。这节课我不讲理论,直接手写一个能用的调度器,把并发数控制住,顺便支持优先级和超时。代码全部原生 JavaScript,无任何依赖,浏览器和 Node 都能跑。

最终效果先看一眼

const scheduler = new Scheduler({ concurrency: 3 });

scheduler.add(() => fetch('/api/log?level=1'), { priority: 1, timeout: 2000 });
scheduler.add(() => fetch('/api/log?level=2'), { priority: 2 });
scheduler.add(() => fetch('/api/log?level=3'), { priority: 3 });

const results = await scheduler.waitAll();

这比用 Promise.all 灵活多了:Promise.all 会把所有请求一次性发出去,浏览器直接卡死。调度器能让同时只有 3 个请求在飞,优先级高的先走,单个任务超过指定时间自动放弃。

任务调度的核心思路

本质就是一个生产者-消费者模型。我们不断往里丢任务,调度器内部维护一个队列,然后启动固定的 worker 去队列里取任务执行。

队列不能是简单 FIFO,因为要支持优先级。这里我省略了复杂的二叉堆,直接用一个数组加 sort,任务量不大的时候性能完全够用。

每个任务的本质是一个返回 Promise 的函数,有没有参数无所谓,只要返回 Promise。这样 fetch、axios、甚至一个 async 函数都能统一处理。

完整代码:先看调度器本体

class Scheduler {
    constructor({ concurrency = 3 } = {}) {
        this.concurrency = Math.max(1, concurrency);
        this._queue = [];
        this._activeCount = 0;
        this._waiters = [];
        this._completedCount = 0;
        this._totalAdded = 0;
    }

    add(task, { priority = 0, timeout = 0 } = {}) {
        if (typeof task !== 'function') {
            throw new TypeError('task must be a function that returns a promise');
        }

        this._totalAdded++;
        const id = this._totalAdded;

        const wrappedTask = () => {
            let timer = null;

            const originalPromise = Promise.resolve().then(task);

            if (timeout > 0) {
                let timeoutResolve;
                let timeoutReject;
                const timeoutPromise = new Promise((resolve, reject) => {
                    timeoutResolve = resolve;
                    timeoutReject = reject;
                });

                timer = setTimeout(() => {
                    timeoutReject(new Error(`Task ${id} timed out after ${timeout}ms`));
                }, timeout);

                const racePromise = Promise.race([originalPromise, timeoutPromise]);

                // 消除未捕获的 rejection 警告
                racePromise.catch(() => {});

                return racePromise.finally(() => clearTimeout(timer));
            }

            return originalPromise;
        };

        this._queue.push({
            id,
            task: wrappedTask,
            priority: Number(priority) || 0,
            createdAt: Date.now(),
        });

        this._queue.sort((a, b) => b.priority - a.priority);

        this._flush();

        return this._createHandle(id);
    }

    _createHandle(id) {
        let settled = false;
        let resolveFn;
        let rejectFn;

        const promise = new Promise((resolve, reject) => {
            resolveFn = resolve;
            rejectFn = reject;
        });

        const checkSettle = () => {
            if (settled) return;
            const idx = this._queue.findIndex(item => item.id === id);
            if (idx === -1) return;

            const item = this._queue[idx];
            if (item._resolved) {
                settled = true;
                resolveFn(item._result);
                this._queue.splice(idx, 1);
            } else if (item._rejected) {
                settled = true;
                rejectFn(item._error);
                this._queue.splice(idx, 1);
            }
        };

        // 注册一个立即检查,因为任务可能已经执行完了
        checkSettle();

        return {
            promise,
            then: (onFulfilled, onRejected) => promise.then(onFulfilled, onRejected),
            catch: (onRejected) => promise.catch(onRejected),
            finally: (onFinally) => promise.finally(onFinally),
        };
    }

    async _flush() {
        while (this._activeCount < this.concurrency && this._queue.length > 0) {
            const next = this._queue.shift();
            if (!next) break;

            this._activeCount++;

            Promise.resolve()
                .then(() => next.task())
                .then(
                    (result) => {
                        next._resolved = true;
                        next._result = result;
                    },
                    (error) => {
                        next._rejected = true;
                        next._error = error;
                    }
                )
                .finally(() => {
                    this._activeCount--;
                    this._completedCount++;
                    this._flush();
                    this._notifyWaiters();
                });
        }
    }

    _notifyWaiters() {
        if (this._waiters.length > 0 && this._activeCount === 0 && this._queue.length === 0) {
            const waiters = this._waiters;
            this._waiters = [];
            waiters.forEach((resolve) => resolve());
        }
    }

    async waitAll() {
        if (this._activeCount === 0 && this._queue.length === 0) {
            return;
        }

        await new Promise((resolve) => {
            this._waiters.push(resolve);
        });
    }

    get pending() {
        return this._queue.length;
    }

    get active() {
        return this._activeCount;
    }

    get completed() {
        return this._completedCount;
    }
}

这个实现比网上常见的调度器多做了三件事:支持超时、支持优先级、每个任务返回一个可 await 的句柄。其中一个关键技巧是用 _flush() 来递归派遣新任务,而不是维护一个常驻 Worker 池。这样代码会简洁很多,同时也避免了空转。

每个 api 都是怎么工作的

add(task, options):把任务包装成带超时的版本,塞进队列,按优先级排序,然后调用 _flush() 尝试启动新任务。返回一个带 then/catch/finally 的对象,这就是一个 Promise 兼容体,可以直接 await。

_flush():只有当活跃任务数小于并发上限时,才会继续从队列头部取出任务来执行。每次任务完成,不管成功还是失败,都会重新调用 _flush(),让队列里的位置空出来。

waitAll():等待所有任务完成。它利用一个 _waiters 数组,如果还有任务在执行,就 push 一个 resolve 函数。当最后一个任务结束时,_notifyWaiters 会触发所有 waitAll 的 Promise resolve。

我故意没有用 async/await 去写 _flush 里面的执行逻辑,而是用 Promise.resolve().then 包了一层。这样即使 task 同步 throw,也能被捕获成 rejected 状态。

真实案例:批量下载并压缩图片

我这边需要把用户上传的十张原图,压缩成三种尺寸。图片接口速度很慢,每张可能要 1s。如果用 Promise.all 同时发 10 个,浏览器网络队列就爆了,而且服务端并发太高容易挂。用调度器就非常自然。

const scheduler = new Scheduler({ concurrency: 2 });

const urls = [
    '/img/photo1.jpg',
    '/img/photo2.jpg',
    '/img/photo3.jpg',
    '/img/photo4.jpg',
    '/img/photo5.jpg',
    '/img/photo6.jpg',
    '/img/photo7.jpg',
];

const downloadTask = (url) => async () => {
    const res = await fetch(url);
    const blob = await res.blob();
    // 这里可以加压缩逻辑
    return blob;
};

const tasks = urls.map((url, index) => {
    return scheduler.add(downloadTask(url), {
        priority: index % 2 === 0 ? 2 : 1, // 偶数的图片优先级高
        timeout: 3000,                      // 3秒超时就放弃
    });
});

const results = await Promise.all(tasks);
console.log('全部完成,共', results.length, '张图片');

因为每个任务都会返回自己的 Promise,所以可以单独 await 某个图片是否成功。如果某个图片超时了,会捕获 Error,不会影响其他任务。这就是我为什么要给每个 add 返回值塞一个 handle 对象的原因。

调试小技巧:在调度器里加日志

如果你想观察并发变化,在 _flush() 的 while 循环里加一行 console.log 就行。

while (this._activeCount < this.concurrency && this._queue.length > 0) {
    const next = this._queue.shift();
    console.log(`[执行] id=${next.id} 活跃=${this._activeCount} 队列=${this._queue.length}`);
    // ...
}

这样你能直观看到当 active 从 0 跳到 2 之后,队列里剩下的任务在等待。等其中一个完成任务,_activeCount 减一,_flush 又会立刻拉一个新的任务进来。

我第一次写完测试,看到日志里最多只有 2 个任务同时执行,瞬间感觉稳了。

关于优先级的一个注意点

我直接用数组 sort 来实现优先级排序。sort 是稳定的,所以当两个任务优先级一样时,先 add 的会排在前面。如果你有非常多的任务(几万条),建议换成一个最小堆,否则每次 add 都全量排序,性能会下降。

另外,优先级并不能中断正在执行的任务,只能影响队列中任务的先后顺序。如果有人正在执行一个优先级低的任务,后来加了一个优先级高的任务,那也只能等当前任务跑完。如果你需要真正抢占式的调度,得配合 AbortController 取消旧任务,这是另一套复杂逻辑,今天先不展开。

超时到底有什么用

超时这功能在网上很多调度器实现里都没有。但我强烈建议加上,因为真实网络环境实在太扯了,一个请求可能卡 10 秒没响应,如果并发池被这种慢请求占满,后面的任务全都堵死。

有了超时之后,超过指定时间直接 reject,然后释放并发槽位,其他任务继续跑。我给的实现里用了 Promise.race 加 setTimeout,注意在 finally 里 clearTimeout,避免内存泄漏。

有个细节:timeoutPromise 可能会先 reject,而 originalPromise 后面才完成。此时 originalPromise 会变成未处理状态吗?不会,因为 Promise.race 已经把 originalPromise 加入订阅了,它内部会处理后来完成的情况。我的代码里还额外加了个 racePromise.catch(() => {}),保证 racePromise 自身如果提前被 reject,不会产生 unhandledrejection 警告。

如果任务本身就是带优先级的队列呢

就比如某些任务需要立即执行,某些可以慢慢排队。你有几种方式:

scheduler.add(criticalTask, { priority: 100 });
scheduler.add(hourlyTask, { priority: 0 });
scheduler.add(backgroundTask, { priority: -10 });

优先级这个字段可以是负数,排序算法会自然处理。这比给每种任务单独建一个队列要方便得多。

最后:集成到现有的请求模块

我目前项目里的所有请求都已经过这层调度器了。换一种说法:它相当于一个本地网关,控制请求的出口速度。这样在并发窗口期,不会因为瞬间的流量打爆后端。

如果你直接用 fetch,也可以这么包装一下:

const http = {
    get(url) {
        return scheduler.add(() => fetch(url), { priority: 2, timeout: 8000 });
    },
    post(url, body) {
        return scheduler.add(() => fetch(url, {
            method: 'POST',
            body: JSON.stringify(body)
        }), { priority: 1, timeout: 10000 });
    }
};

然后项目里的代码就从直接 fetch 换成了 http.get 和 http.post。好处是全局都能统一控制并发和超时,不用每个页面各自实现一遍。

顺便说一下,这个调度器对 Node.js 同样有效。在 Node 里做爬虫,控制并发也是必须的。把 fetch 换成 axios 就行。

你能改造成什么样

我的实现里没有用到 generator、没有用到 worker_threads,就是最朴素的 Promise 队列。如果你觉得功能不够,可以加上任务的取消、重试、进度回调、甚至暂停/恢复整个队列。

比如加一个重试机制,在任务执行失败时重新入队,最大重试 3 次。这个方法可以放在 _flush 里面,当 .then 的 reject 分支里判断重试次数。但要注意别造成无限循环,给个上限就好。

整个调度器不过一百多行,但能解决的是前端里最头疼的并发控制问题。先写一个能用的,再慢慢加功能。代码本来就是反复打磨出来的。

希望这篇能让你少走弯路。如果你按我代码跑,遇到任何问题,直接来看我的逻辑,肯定能调通。毕竟我已经在浏览器和 Node 里跑过无数遍了。

JavaScript并发控制实战:手写带优先级与超时的任务调度器
收藏 (0) 打赏

感谢您的支持,我会继续努力的!

打开微信/支付宝扫一扫,即可进行扫码打赏哦,分享从这里开始,精彩与您同在
点赞 (0)

版权声明:
本站资源有的来自互联网收集整理,本站纯免费分享提供学习使用,如果侵犯了您的合法权益,请发送邮件1506151422@qq.com联系,将会及时下架删除。
本站资源仅供研究、学习交流之用,免费开源项目不代表完全可商用,若商业用途请先咨询开发企业能否商用,否则产生的一切后果将由下载用户自行承担。
原创板块未经允许不得转载,否则将追究法律责任。

淘吗网 javascript JavaScript并发控制实战:手写带优先级与超时的任务调度器 https://www.taomawang.com/web/javascript/2592.html

常见问题

相关文章

猜你喜欢
发表评论
暂无评论
官方客服团队

为您解决烦忧 - 24小时在线 专业服务