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

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;

额外优化建议

  1. Flask端改用生产级服务器:默认的Werkzeug开发服务器并发能力有限,建议使用Gunicorn或uWSGI提升性能,示例命令:

    gunicorn --workers=4 --threads=8 --bind=0.0.0.0:5000 app:app
    

    根据服务器硬件资源调整workers和threads参数。

  2. 添加流状态监控:可以增加日志或监控逻辑,跟踪每个流的订阅者数量、连接状态,方便后续排查问题。

  3. 处理异常边界:代码中已加入写入失败时移除订阅者的逻辑,避免无效连接占用资源。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 15:47:13