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

Node.js新手开发Webhook:MQTT异步等消息后响应HTTP请求

解决Node.js Webhook + MQTT 关联响应的问题

嘿,作为Node.js新手能尝试Webhook和MQTT的集成,已经很赞啦!我看了你的代码片段,发现几个核心问题需要调整,才能实现你想要的「接收HTTP POST → 发送MQTT请求 → 等待关联MQTT消息 → 返回HTTP响应」的完整流程。下面一步步给你拆解和优化:

你的代码目前存在的问题

  • 全局变量引发并发冲突:你用了全局的array和count来存储MQTT消息,要是同时有多个用户发POST请求,这些请求的消息会混在一起,根本分不清哪条消息属于哪个请求。
  • 订阅时机不合理:在POST请求里订阅inTopic会导致重复订阅,而且所有请求都会收到这个topic的所有消息,无法关联到特定的请求。
  • 缺少响应触发逻辑:代码里没有在收到MQTT消息后主动给HTTP请求返回响应,会导致请求一直挂着超时。

优化后的完整实现方案

我们可以用请求唯一ID来关联HTTP请求和MQTT消息,再用Map来暂存等待中的请求,这样就能精准匹配每个请求对应的MQTT响应:

const express = require('express');
const mqtt = require('mqtt');
const { v4: uuidv4 } = require('uuid'); // 用来生成唯一请求ID

const app = express();
app.use(express.json()); // 解析POST请求的JSON body

// MQTT配置
const MQTTServer = 'mqtt://your-mqtt-broker-url';
const client = mqtt.connect(MQTTServer);

// 存储等待中的请求:key是requestId,value是resolve函数
const pendingRequests = new Map();

// MQTT连接成功后订阅主题(只订阅一次)
client.on('connect', () => {
  console.log('MQTT客户端已连接');
  client.subscribe('inTopic', (err) => {
    if (err) console.error('订阅失败:', err);
  });
});

// 处理收到的MQTT消息
client.on('message', (topic, messageBuffer) => {
  try {
    const message = JSON.parse(messageBuffer.toString());
    const { requestId, data } = message;

    // 找到对应的等待请求
    if (pendingRequests.has(requestId)) {
      const resolve = pendingRequests.get(requestId);
      // 返回MQTT消息作为HTTP响应
      resolve({ status: 'success', data });
      // 清理已完成的请求
      pendingRequests.delete(requestId);
    }
  } catch (err) {
    console.error('解析MQTT消息失败:', err);
  }
});

// 处理Webhook的POST请求
app.post('/test', async (request, response) => {
  try {
    // 生成唯一请求ID,用来关联MQTT消息
    const requestId = uuidv4();

    // 构造要发送的MQTT消息,带上requestId
    const mqttMessage = JSON.stringify({
      requestId,
      payload: request.body // 把POST请求的内容传给MQTT
    });

    // 发布MQTT消息
    client.publish('outTopic', mqttMessage, (err) => {
      if (err) {
        response.status(500).json({ status: 'error', message: 'MQTT发布失败' });
        return;
      }
    });

    // 等待MQTT响应,设置10秒超时(避免请求一直挂着)
    const timeoutPromise = new Promise((_, reject) => {
      setTimeout(() => {
        reject(new Error('等待MQTT响应超时'));
        pendingRequests.delete(requestId);
      }, 10000);
    });

    // 等待MQTT消息或者超时
    const result = await Promise.race([
      new Promise(resolve => pendingRequests.set(requestId, resolve)),
      timeoutPromise
    ]);

    // 返回响应给客户端
    response.json(result);
  } catch (err) {
    response.status(500).json({ status: 'error', message: err.message });
  }
});

// 启动HTTP服务
const PORT = 3000;
app.listen(PORT, () => {
  console.log(`Webhook服务运行在 http://localhost:${PORT}`);
});

关键优化点说明

  • 请求唯一标识:用uuid生成requestId,让每个HTTP请求和对应的MQTT消息一一对应,解决并发冲突问题。
  • 提前订阅MQTT主题:在MQTT客户端连接成功后只订阅一次inTopic,避免重复订阅。
  • 暂存等待请求:用Map存储每个请求的resolve函数,收到匹配的MQTT消息时触发响应。
  • 超时处理:给等待MQTT消息设置超时时间,防止请求无限期挂起。

额外注意事项

  • 确保你的MQTT服务端返回的消息里包含requestId,这样才能匹配到对应的HTTP请求。
  • 要处理MQTT连接断开的情况,可以添加client.on('reconnect')和client.on('error')事件监听,保证服务稳定性。
  • 如果需要等待多条MQTT消息,可以修改逻辑为收集指定数量的消息后再返回响应。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 08:11:01