Node.js与MQTT集成报错:Cannot set headers after they are sent to the client
问题描述
我开发了冰箱控制器的Node.js后端,集成MQTT协议实现设备交互。其中ShowCurrentTemperature接口通过订阅MQTT的temperature主题获取并返回温度数据,首次调用正常,但后续调用会抛出Cannot set headers after they are sent to the client错误。尝试过添加return语句等方法,均未解决问题,恳请协助排查。
原控制器代码:
const fs = require('fs'); const mqtt = require('mqtt'); const transporter = require('../params/mail') const winston = require('../params/log'); const User = require("../models/User"); const { cli } = require('winston/lib/winston/config'); exports.OpenTheCase = async (req, res) => {}; exports.AddCard = async (req, res) => {}; exports.ShowCurrentTemperature = async (req, res) => {}; exports.ShowCurrentHumidity = async (req, res) => {}; exports.SetAlarm = async (req, res) => {}; exports.AlarmIsOn = async (req, res) => {}; const options = { clientId: 'backendserver1032', key: fs.readFileSync('./certs/mqtt_cert/client.key'), cert: fs.readFileSync('./certs/mqtt_cert/client.crt'), ca: [ fs.readFileSync('./certs/mqtt_cert/ca.crt') ] } const client = mqtt.connect('mqtts://localhost:8883', options); exports.OpenTheCase = async (req, res) => { try { client.publish('RFID', 'RFID_OPEN'); res.status(200).json({ 'case':"opened" }); } catch(e){ res.status(200).json({ 'state':"something went wrong" }); } } exports.AddCard = async (req, res) => { try { client.publish('RFID', 'RFID_ADD'); res.status(200).json({ 'card':"will be added" }); } catch(e){ res.status(200).json({ 'state':"something went wrong" }); } } exports.ShowCurrentTemperature = async (req, res) => { try { client.subscribe('temperature'); client.on('message', (topic, message, packet) => { res.status(200).json({ 'temperature': message.toString('ascii') }) client.unsubscribe('temperature') }) } catch(e){ res.status(200).json({ 'state':"something went wrong" }); } return } exports.ShowCurrentHumidity = async (req, res) => { try { client.subscribe('humidity'); client.on('message', (topic, message) => { res.status(200).json({"temperature": message.toString('ascii')}); client.unsubscribe('humidity') }); } catch(e){ res.status(200).json({ 'state':"something went wrong" }); } return }
路由配置:
router.get("/frigo/Temperature",auth.verifyToken, frigoController.ShowCurrentTemperature)
问题根源
- 每次调用
ShowCurrentTemperature或ShowCurrentHumidity接口时,都会执行client.on('message', ...),这会新增一个事件监听器,而非覆盖原有监听器。 - 当后续MQTT消息到达时,所有之前绑定的监听器都会被触发,每个监听器都尝试向对应的已结束请求的响应对象(
res)发送数据,从而抛出“Cannot set headers after they are sent to the client”错误。 - 频繁调用
client.subscribe和client.unsubscribe会导致MQTT订阅逻辑混乱,进一步加剧问题。
解决方案
提供两种可行方案,按需选择:
方案1:使用一次性事件监听器(快速修复)
通过手动管理事件监听器的绑定与移除,确保每个请求只对应一个监听器,避免重复触发:
exports.ShowCurrentTemperature = async (req, res) => { try { // 定义消息处理回调 const handleMessage = (topic, message) => { if (topic !== 'temperature') return; // 返回温度数据 res.status(200).json({ 'temperature': message.toString('ascii') }); // 移除监听器,避免内存泄漏 client.off('message', handleMessage); // 取消订阅当前主题 client.unsubscribe('temperature'); }; // 绑定监听器 client.on('message', handleMessage); // 订阅主题,处理订阅失败情况 client.subscribe('temperature', (err) => { if (err) { res.status(500).json({ 'state': '订阅温度主题失败' }); client.off('message', handleMessage); } }); // 添加超时处理,防止请求无限挂起 const timeout = setTimeout(() => { res.status(504).json({ 'state': '获取温度超时' }); client.off('message', handleMessage); client.unsubscribe('temperature'); }, 10000); // 请求完成后清除超时定时器 res.on('finish', () => clearTimeout(timeout)); } catch (e) { res.status(500).json({ 'state': '服务器内部错误' }); } } exports.ShowCurrentHumidity = async (req, res) => { try { const handleMessage = (topic, message) => { if (topic !== 'humidity') return; res.status(200).json({ 'humidity': message.toString('ascii') }); client.off('message', handleMessage); client.unsubscribe('humidity'); }; client.on('message', handleMessage); client.subscribe('humidity', (err) => { if (err) { res.status(500).json({ 'state': '订阅湿度主题失败' }); client.off('message', handleMessage); } }); const timeout = setTimeout(() => { res.status(504).json({ 'state': '获取湿度超时' }); client.off('message', handleMessage); client.unsubscribe('humidity'); }, 10000); res.on('finish', () => clearTimeout(timeout)); } catch (e) { res.status(500).json({ 'state': '服务器内部错误' }); } }
方案2:缓存最新温湿度数据(推荐,性能更优)
在MQTT客户端连接后长期订阅温湿度主题,缓存最新数据,接口直接返回缓存值,避免频繁订阅/取消订阅:
// 初始化缓存变量 let latestTemperature = null; let latestHumidity = null; // 客户端连接成功后订阅主题 client.on('connect', () => { console.log('MQTT客户端已成功连接'); client.subscribe(['temperature', 'humidity'], (err) => { if (err) console.error('订阅温湿度主题失败:', err); }); }); // 监听MQTT消息,更新缓存 client.on('message', (topic, message) => { const payload = message.toString('ascii'); switch(topic) { case 'temperature': latestTemperature = payload; break; case 'humidity': latestHumidity = payload; break; } }); // 温度接口直接返回缓存数据 exports.ShowCurrentTemperature = async (req, res) => { try { if (!latestTemperature) { return res.status(404).json({ 'state': '暂无温度数据' }); } res.status(200).json({ 'temperature': latestTemperature }); } catch (e) { res.status(500).json({ 'state': '服务器内部错误' }); } } // 湿度接口直接返回缓存数据 exports.ShowCurrentHumidity = async (req, res) => { try { if (!latestHumidity) { return res.status(404).json({ 'state': '暂无湿度数据' }); } res.status(200).json({ 'humidity': latestHumidity }); } catch (e) { res.status(500).json({ 'state': '服务器内部错误' }); } }
内容的提问来源于stack exchange,提问作者Arthur Secondaire
相关产品推荐
相关产品推荐

