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

基于cacheKey的最大并发限制并发队列优化咨询

基于CacheKey的并发队列优化方案

问题背景与现有实现缺陷

我已实现一个基于cacheKey的并发队列,用于控制每个用户(以cacheKey区分)最多4个并发调用外部API,但现有实现存在以下问题:

  • 采用4个独立FIFO队列均匀分配任务,因无法预知请求耗时,可能导致整体运行时间不必要拉长;
  • 无法缓存API返回结果,相同任务无法复用结果;
  • 队列控制逻辑冗余,通过额外Promise实现等待,易出错。

外部API限制每个用户最多4个并发请求,应用模拟多用户时并发量可达20,超限制会触发问题。Promise.all分块批量处理效率极低,理想状态是FIFO队列,一个任务完成后立即启动下一个任务,且结果能直接返回给依赖组件。

当前实现代码

class QueuingService {
    /** ~FIFO queue */
    private static queue: Record<string, { promise: Promise<Function>, uuid: string }[]> = {};
    private static MAX_CONCURRENCY = 4

    /** cacheKey is used as cachekey for grouping in the queue */
    public static async addToQueueAndExecute(func: Function, cacheKey: string) {
        let resolver: Function | null = null;
        let promise = new Promise<Function>(resolve => {
            resolver = resolve;
        });
        //in the real world, this is a real key created by nanoId
        let uuid = `${Math.random()}`;


        Array.isArray(this.queue[cacheKey])
            ? this.queue[cacheKey].push({ promise: promise, uuid: uuid })
            : this.queue[cacheKey] = [{ promise: promise, uuid: uuid }];

        //queue in all calls, until MAX_CONCURRENCY is reached. After that, slice the first entry and await the promise
        if (this.queue[cacheKey].length > this.MAX_CONCURRENCY) {
            let queuePromise = this.queue[cacheKey].shift();
            if (queuePromise){
                await queuePromise.promise;
            }
        }
        //console.log("elements in queue:", this.queue[cacheKey].length, cacheKey)

        //technically this wrapping is not necessary, but it makes to code more readable imho
        let result = async () => {
            let res = await func();
            if (resolver) {
                //resolve and clean up
                resolver();
                this.queue[cacheKey] = this.queue[cacheKey].filter(elem => elem.uuid !== uuid);
            }
            //console.log(res, cacheKey, "finshed after", new Date().getTime() - enteredAt.getTime(),"ms", "entered at:", enteredAt)

            return res;
        }

        return await result();
    }
}

async function sleep(ms:number){
    return await new Promise(resolve=>window.setTimeout(resolve,ms))
}

async function caller(){
    /* //just for testing with fewer entries
    for(let i=0;i<3;i++){
        let groupKey = Math.random() <0.5? "foo":"foo"
        QueuingService.addToQueueAndExecute(async ()=>{
            await sleep(Math.floor(4000-Math.random()*2000));
            console.log(i, new Date().getSeconds(), new Date().getMilliseconds(),groupKey)
            return Math.random()
        },groupKey)
    }
    */
    
    for(let i=0;i<20;i++){
        let groupKey = Math.random() <0.5? "foo":"foo";
        let startedAt = new Date();
        QueuingService.addToQueueAndExecute(async ()=>{
            await sleep(Math.floor(4000-Math.random()*2000));
            console.log(i, new Date().getTime()-startedAt.getTime(),groupKey)
            return Math.random()
        },groupKey)
    }
    
}
caller()

优化建议与实现方案

1. 核心优化思路

  • 每个cacheKey维护单个任务队列+当前并发数计数器,替代多队列分配,保证一有空闲就启动下一个任务;
  • 简化Promise控制逻辑,直接绑定任务执行Promise,避免冗余的resolve操作;
  • 新增结果缓存,支持复用已完成任务的结果;
  • 提升类型安全性,用泛型支持不同任务返回类型。

2. 优化后的代码实现

import { nanoid } from 'nanoid'; // 实际项目中使用nanoid生成唯一ID

class QueuingService {
    // 每个cacheKey的任务队列:存储待执行的任务函数
    private static taskQueues: Record<string, (() => Promise<any>)[]> = {};
    // 每个cacheKey的当前并发数
    private static activeCounts: Record<string, number> = {};
    // 结果缓存:key为cacheKey+任务唯一标识,value为任务结果
    private static resultCache: Record<string, any> = {};
    private static readonly MAX_CONCURRENCY = 4;

    /**
     * 添加任务到队列并执行
     * @param func 待执行的API调用函数(返回Promise)
     * @param cacheKey 用户标识分组键
     * @param taskId 任务唯一标识(用于缓存,不传则自动生成)
     */
    public static async addToQueueAndExecute<T>(
        func: () => Promise<T>,
        cacheKey: string,
        taskId?: string
    ): Promise<T> {
        const finalTaskId = taskId || nanoid();
        const cacheResultKey = `${cacheKey}:${finalTaskId}`;

        // 先检查缓存,存在则直接返回
        if (this.resultCache[cacheResultKey]) {
            return this.resultCache[cacheResultKey];
        }

        // 初始化队列和计数器
        if (!this.taskQueues[cacheKey]) {
            this.taskQueues[cacheKey] = [];
            this.activeCounts[cacheKey] = 0;
        }

        return new Promise<T>((resolve, reject) => {
            // 将任务包装成带resolve/reject的函数加入队列
            this.taskQueues[cacheKey].push(async () => {
                try {
                    const result = await func();
                    // 缓存结果
                    this.resultCache[cacheResultKey] = result;
                    resolve(result);
                } catch (err) {
                    reject(err);
                } finally {
                    // 并发数减1,执行下一个任务
                    this.activeCounts[cacheKey]--;
                    this.processNextTask(cacheKey);
                }
            });

            // 尝试启动任务
            this.processNextTask(cacheKey);
        });
    }

    /**
     * 处理队列中的下一个任务
     */
    private static processNextTask(cacheKey: string) {
        const queue = this.taskQueues[cacheKey];
        const activeCount = this.activeCounts[cacheKey];

        // 队列有任务且并发数未达上限时启动任务
        if (queue.length > 0 && activeCount < this.MAX_CONCURRENCY) {
            this.activeCounts[cacheKey]++;
            const nextTask = queue.shift();
            nextTask?.();
        }
    }

    /**
     * 清除指定cacheKey的队列和缓存
     */
    public static clearCache(cacheKey: string) {
        delete this.taskQueues[cacheKey];
        delete this.activeCounts[cacheKey];
        // 清除该cacheKey下的所有缓存项
        Object.keys(this.resultCache).forEach(key => {
            if (key.startsWith(`${cacheKey}:`)) {
                delete this.resultCache[key];
            }
        });
    }
}

// 测试代码
async function sleep(ms: number) {
    return await new Promise(resolve => setTimeout(resolve, ms));
}

async function caller() {
    for (let i = 0; i < 20; i++) {
        const groupKey = Math.random() < 0.5 ? "foo" : "bar";
        const startedAt = new Date();
        QueuingService.addToQueueAndExecute(async () => {
            await sleep(Math.floor(4000 - Math.random() * 2000));
            const duration = new Date().getTime() - startedAt.getTime();
            console.log(`任务${i}完成,耗时${duration}ms,分组:${groupKey}`);
            return Math.random();
        }, groupKey, `task-${i}`)
        .catch(err => console.error(`任务${i}失败:`, err));
    }
}

caller();

3. 优化点说明

  • 并发效率提升:单个队列+并发计数器的模式,保证只要有并发额度就立即启动下一个任务,避免因任务耗时不均导致的资源闲置;
  • 结果缓存:通过taskId和cacheKey组合的键缓存结果,相同任务可直接复用;
  • 逻辑简化:移除冗余的Promise等待逻辑,任务完成后自动触发下一个任务,代码更清晰;
  • 类型安全:泛型支持不同任务返回类型,符合TypeScript最佳实践;
  • 可维护性:新增clearCache方法,支持清理指定用户的队列和缓存,避免内存泄漏。

内容的提问来源于stack exchange,提问作者user13897241

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 09:50:17