Express.js仅支持6个客户端连接,如何取消连接限制?
解决Express中转Flask图像流的连接数限制问题
问题根源
你当前的架构是每个客户端请求都会新建一个到Flask的长连接,而Flask端(尤其是默认的Werkzeug开发服务器)或Node.js HTTP客户端的连接池限制了并发连接数,导致最多仅支持6个客户端。修改server.maxConnections无效是因为该参数控制的是Express自身的总连接数,瓶颈实际在于Express到Flask的连接复用率过低。
核心解决方案:复用Flask连接,广播流给多客户端
不为每个客户端单独建立Flask连接,而是为每个cctvId+port维护一个共享的Flask流,将收到的图像帧广播给所有订阅该流的客户端。这样无论多少客户端,仅占用一个Flask连接,彻底突破连接数限制。
修改后的stream.js代码
process.env.NODE_TLS_REJECT_UNAUTHORIZED = '0'; const express = require('express'); const router = express.Router(); const http = require('http'); const https = require('https'); // 维护每个cctv+port的流实例和订阅者列表 const streamSubscribers = new Map(); // 创建全局Agent,确保客户端连接池无限制 const httpAgent = new http.Agent({ maxSockets: Infinity }); const httpsAgent = new https.Agent({ maxSockets: Infinity }); router.get('/api/cctv/:cctvId*?', (req, res) => { const cctvId = req.params.cctvId; const port = req.query.port; const streamKey = `${cctvId}_${port}`; // 设置响应头 res.writeHead(200, { 'Content-Type': 'multipart/x-mixed-replace; boundary=frame', 'Connection': 'close', 'Transfer-Encoding': 'chunked', }); // 客户端断开时的清理逻辑 const cleanup = () => { const streamEntry = streamSubscribers.get(streamKey); if (streamEntry) { streamEntry.subscribers.delete(res); // 无订阅者时关闭Flask连接 if (streamEntry.subscribers.size === 0) { streamEntry.stream.abort(); streamSubscribers.delete(streamKey); console.log(`Closed Flask stream for ${streamKey}`); } } res.end(); }; res.on('close', cleanup); res.on('error', cleanup); // 检查是否已有共享的Flask流 let streamEntry = streamSubscribers.get(streamKey); if (!streamEntry) { const isHttps = false; // 根据实际Flask协议调整 const client = isHttps ? https : http; const agent = isHttps ? httpsAgent : httpAgent; // 建立Flask连接 const flaskRequest = client.get({ path: `/image/infer/play/${cctvId}`, hostname: 'localhost', port: port, agent: agent }, (flaskRes) => { flaskRes.on('data', (chunk) => { // 将图像帧广播给所有订阅者 const currentEntry = streamSubscribers.get(streamKey); if (currentEntry) { currentEntry.subscribers.forEach(subscriberRes => { try { subscriberRes.write(chunk); } catch (err) { // 写入失败,移除异常订阅者 currentEntry.subscribers.delete(subscriberRes); } }); } }); flaskRes.on('end', () => { // Flask流结束时,关闭所有订阅者连接 const currentEntry = streamSubscribers.get(streamKey); if (currentEntry) { currentEntry.subscribers.forEach(subscriberRes => subscriberRes.end()); streamSubscribers.delete(streamKey); } }); }); flaskRequest.on('error', (err) => { console.log(`Failed to connect Flask stream ${streamKey}:`, err); const currentEntry = streamSubscribers.get(streamKey); if (currentEntry) { currentEntry.subscribers.forEach(subscriberRes => subscriberRes.end()); streamSubscribers.delete(streamKey); } }); // 初始化订阅者列表 streamEntry = { stream: flaskRequest, subscribers: new Set([res]) }; streamSubscribers.set(streamKey, streamEntry); console.log(`Started Flask stream for ${streamKey}`); } else { // 已有流,添加当前客户端到订阅者列表 streamEntry.subscribers.add(res); } }); module.exports = router;
额外优化建议
Flask端改用生产级服务器:默认的Werkzeug开发服务器并发能力有限,建议使用Gunicorn或uWSGI提升性能,示例命令:
gunicorn --workers=4 --threads=8 --bind=0.0.0.0:5000 app:app根据服务器硬件资源调整
workers和threads参数。添加流状态监控:可以增加日志或监控逻辑,跟踪每个流的订阅者数量、连接状态,方便后续排查问题。
处理异常边界:代码中已加入写入失败时移除订阅者的逻辑,避免无效连接占用资源。
内容的提问来源于stack exchange,提问作者Hyeonjun Jeong
相关产品推荐
相关产品推荐

