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

Node.js集群中MQTT.js连接异常关闭问题求助

解决Node.js集群中MQTT客户端连接异常问题

在Node.js集群环境中引入mqtt.js文件后,MQTT客户端无法正常运行,连接被关闭并报错。以下是涉及的代码文件:

集群文件

const cluster = require('cluster');
const os = require ("os");

const cpuCount = os.cpus().length;
console.log(`The total number of CPUs is ${cpuCount}`);

console.log(`Primary pid=${process.pid}`);
cluster.setupPrimary({
  exec: __dirname + "/app.js",
});

for (let i = 0; i < cpuCount; i++) {
  cluster.fork();
}
cluster.on("exit", (worker, code, signal) => {
  console.log(`worker ${worker.process.pid} has been killed`);
  console.log("Starting another worker");
  cluster.fork();
});

mqtt.js文件

//------- MQTTT --------//

const mqtt = require("mqtt");
const queryFunction = require("./utilities/queryFunction")
const client = mqtt.connect("mqtt://greenenergyhomes.uniupo.click", {clientId: "Archiver"});

const topic = ["event/greenenergy/+", "desc/greenenergy/+"]

client.on('connect',() => {
  console.log("connessione con il broker MQTT stabilita");
  client.subscribe(topic, console.log);
  
});

client.on('message',(topic, message) => {
    console.log(`message: ${message}, topic: ${topic}`);
    const msgJSON = JSON.parse(message);

    const params = topic.split("/");
    console.log(params)
    switch(params[0]){
      case 'event':
        queryFunction.updateVal(msgJSON, params[3]);
        break;
      case 'desc':
        console.log("desc in MQTT")
        queryFunction.setVal(msgJSON, params[2])
        break
    } 

  });

module.exports;

问题原因及解决办法

  • 重复ClientID触发Broker连接拒绝
    MQTT Broker不允许相同clientId的多个客户端同时连接,集群中每个Worker进程都会加载mqtt.js并创建clientId: "Archiver"的客户端,新连接会踢掉旧连接,导致循环断开。
    解决:给每个Worker的ClientID添加唯一标识,比如进程ID:

    const client = mqtt.connect("mqtt://greenenergyhomes.uniupo.click", {
      clientId: `Archiver-${process.pid}`
    });
    
  • 缺少错误监听无法排查问题
    当前代码未监听MQTT客户端的error事件,无法获取具体报错信息。
    解决:在mqtt.js中添加错误监听:

    client.on('error', (err) => {
      console.error(`MQTT客户端错误 (进程ID: ${process.pid}):`, err);
    });
    
  • 可选:优化集群环境下的MQTT连接架构
    如果不需要每个Worker独立处理MQTT消息,可以改为在主进程创建单个MQTT客户端,通过IPC通信将消息分发给Worker处理,减少Broker连接压力:

    1. 主进程中引入mqtt.js(修改为仅创建客户端并订阅消息)
    2. 主进程收到MQTT消息后,通过worker.send()发送给Worker
    3. Worker监听message事件处理业务逻辑

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 13:21:26