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
相关产品推荐
相关产品推荐

