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

Node.js中EventEmitter结合MQTT事件处理器发布失效问题排查

问题分析与解决方案

核心问题:异步操作未完成就终止进程

你的独立测试脚本在触发setSysDetails和close mqtt事件后立即调用process.exit(1),但MQTT的publish是异步操作——它仅将消息加入发送队列,不会立即完成。进程提前终止会导致事件循环无法处理未完成的MQTT发送任务,最终消息无法到达代理。

具体修复步骤

1. 避免提前终止进程,等待异步操作完成

修改testEvent.js逻辑,不要直接调用process.exit(1),而是通过事件监听等待MQTT操作完成后再退出:

// testEvent.js 修改后
const emitter = require("./config/EventBus.js");

emitter.emit("lcl_main_start");

const dao = new AppDAO(constants.dbFile);
const siteDao = new SiteDAO(dao);

// 监听MQTT关闭完成事件,再退出进程
emitter.on("mqtt_closed", () => {
    console.log(true);
    process.exit(0); // 使用0表示正常退出
});

emitter.on("eventBus_init", () => {
    try {
      let tick = new dayjs();
      siteDao.getActive().then(localSite => {
        emitter.emit("setSysDetails", "buildUpTime", (dayjs() - tick));
        emitter.emit("close mqtt"); // 触发关闭,等待完成后退出
      }).catch(err => {
          console.log(err);
          emitter.emit("close mqtt");
      });
    } catch (err) {
        console.log(err);
        emitter.emit("close mqtt");
    }
});

2. 修正MQTT关闭逻辑,确保消息发送完成后再断开连接

在EventBus.js的close mqtt事件处理器中,将client.end()移至publish的回调函数内,保证消息发送完成后再关闭连接:

// EventBus.js 修改后的close mqtt处理器
emitter.on('close mqtt', () => {
    client.publish("STATUS/systemDetails/connected", JSON.stringify(0), {qos: 0, retain: true}, (err) => {
        if (err) logger.error("MQTT publish error [STATUS/eventBus/connected] : " + err.message);
        // 消息发送完成后再关闭连接
        client.end(() => {
            logger.info('MQTT connection closed');
            emitter.emit("mqtt_closed"); // 通知测试脚本可以退出
        });
    });
});

3. 修复无效的reject调用

setSysDetails处理器中的reject(err)没有对应的Promise上下文,会导致未捕获错误,替换为日志输出:

// EventBus.js 修改后的setSysDetails处理器
emitter.on("setSysDetails", function (param, val) {
    let ts = new dayjs().format("YYYY-MM-DDTHH:mm:ssZ[Z]");
    if (client.connected) {
        client.publish("STATUS/systemDetails/" + param, JSON.stringify({val: val, timestamp: ts}), {qos: 0, retain: true}, (err) => {
            if (err) {
                logger.error(`MQTT publish error [STATUS/systemDetails/${param}] : ${err.message}`);
            }
        });
    } else {
        logger.warn("MQTT client not connected, skipping publish");
    }
});

4. 修正MQTT配置中的类型错误

mqttOptions中的reconnectPeriod被设置为字符串'1000',但MQTT库期望数字类型(毫秒数),改为:

let mqttOptions = {
    // ...其他配置
    reconnectPeriod: 1000, // 去掉引号,改为数字
    // ...其他配置
};

为什么在Express中正常工作?

Express应用会保持事件循环持续运行(等待HTTP请求),所以MQTT的异步publish操作有足够时间完成。而独立脚本默认在同步代码执行完后就终止事件循环,导致异步任务被中断。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 07:05:03