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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 22:01:01