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

NodeJS Express API票务/队列系统 多EAP GRPC并发上限管控实现咨询

两层并发限制队列调度实现方案

整体架构设计

基于你提到的两层Bull队列思路+事件通知机制即可完全满足需求,无需引入RabbitMQ这类额外中间件,整体逻辑如下:

  • 第一层:全局Express总队列,设置concurrency为你定义的总并发上限(比如100),控制所有往外发的EAP请求总数量
  • 第二层:每个EAP对应独立的并发控制器,设置各自的并发上限(比如EAP1设30),控制单EAP的请求并发数
  • 新增全局事件触发器,负责在任意EAP有空位释放时通知总调度器重新调度积压的对应EAP请求

核心模块实现

1. EAP 并发控制器

每个EAP对应一个实例,负责管理当前并发数、空位锁定、释放、空位事件通知:

const EventEmitter = require('events');
class EAPHandler extends EventEmitter {
  constructor(eapName, maxConcurrency) {
    super();
    this.eapName = eapName;
    this.maxConcurrency = maxConcurrency;
    this.currentConcurrency = 0;
  }
  // 申请空位
  acquire() {
    if (this.currentConcurrency < this.maxConcurrency) {
      this.currentConcurrency++;
      return true;
    }
    return false;
  }
  // 释放空位
  release() {
    this.currentConcurrency--;
    // 释放后触发空位通知事件
    this.emit('available', this.eapName);
  }
}
// 初始化所有EAP的控制器,按实际业务调整各EAP的并发上限
const eapHandlers = {
  EAP1: new EAPHandler('EAP1', 30),
  EAP2: new EAPHandler('EAP2', 50),
  EAP3: new EAPHandler('EAP3', 40)
}

2. TicketHandler 全局调度器

负责票据管理、总并发控制、调度逻辑:

const Queue = require('bull');
// 全局总队列,redis配置按实际环境调整
const globalQueue = new Queue('express-eap-global', {
  redis: { host: '127.0.0.1', port: 6379 }
});
// Express全局转发并发上限,按实际需求调整
const GLOBAL_CONCURRENCY = 100;

class TicketHandler {
  static async getTicket(variousInfo, callback) {
    // 把请求信息、回调、目标EAP存入队列,设置无限重试避免请求丢失
    await globalQueue.add({
      variousInfo,
      targetEAP: variousInfo.targetEAP, // 提前确定请求需要调用的EAP
      callback
    }, { attempts: Infinity, backoff: 0 });
  }
  // 初始化调度逻辑
  static init() {
    // 监听所有EAP的空位事件,收到后触发总队列重新扫描待处理任务
    Object.values(eapHandlers).forEach(handler => {
      handler.on('available', () => processQueue());
    });
    // 队列处理公共逻辑
    const processQueue = async (job, done) => {
      const { targetEAP, callback } = job.data;
      const eapHandler = eapHandlers[targetEAP];
      // 申请对应EAP的空位
      if (eapHandler.acquire()) {
        // 申请成功触发回调
        callback();
        done();
        return;
      }
      // 申请失败把任务放回队列末尾,等待下一次调度
      await job.retry();
      done(new Error('EAP concurrency full, retry later'));
    };
    // 启动全局队列处理
    globalQueue.process(GLOBAL_CONCURRENCY, processQueue);
  }
}

3. Express 中间件集成

// 服务启动时先初始化TicketHandler
TicketHandler.init();

// 队列中间件
app.use(async (req, res, next) => {
  // 按实际业务逻辑判断当前请求需要调用的EAP,存入请求信息
  req.variousInfo = {
    targetEAP: 'EAP1', 
    // 其他业务所需的基础信息
  };
  // 申请票据,回调触发后走后续业务逻辑
  await TicketHandler.getTicket(req.variousInfo, () => {
    // 业务逻辑处理完成后,自动释放对应EAP的空位
    res.on('finish', () => {
      eapHandlers[req.variousInfo.targetEAP].release();
    });
    res.on('close', () => {
      eapHandlers[req.variousInfo.targetEAP].release();
    });
    next();
  });
})

关键逻辑说明

  • 重试触发机制:不需要额外定时扫描,任意EAP释放空位时都会触发available事件,直接通知全局队列重新扫描所有待调度的票据,优先处理当前有空位的EAP对应的请求,避免无效重试
  • 调度优先级:全局队列遵循FIFO规则,遇到当前EAP满额的请求会放到队列末尾,优先处理后续匹配到有空位EAP的请求,不会阻塞整体队列
  • 并发控制完全符合要求:第一层全局队列的concurrency参数严格控制总并发数不超过上限,第二层每个EAPHandler的currentConcurrency严格控制单EAP的并发数

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.03 07:45:00