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

