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

Node.js MQTT回调代码改造为Async/Await的响应返回问题求助

改造方案

核心逻辑是将原回调模式替换为手动封装Promise,把Promise的resolve/reject方法挂载到waitingMessageIds的对应记录上,收到响应或触发超时的时候分别调用这两个方法即可完成异步结果返回。


1. 改造publishMessage方法

不再返回messageId,改为返回Promise,将resolve/reject存入等待队列:

async publishMessage(method, namespace, payload) {
  this.clientResponseTopic = `/app/${this.userId}-${appId}/subscribe`;
  const messageId = crypto.createHash('md5').update(generateRandomString(16)).digest('hex');
  const timestamp = Math.round(new Date().getTime() / 1000);
  const signature = crypto.createHash('md5').update(messageId + this.key + timestamp).digest('hex');

  const data = {
    header: {
      from: this.clientResponseTopic,
      messageId, 
      method, 
      namespace, 
      payloadVersion: 1,
      sign: signature, 
      timestamp,
    },
    payload,
  };

  // 捕获MQTT发布错误直接抛出
  try {
    await this.client.publish(`/appliance/${this.uuid}/subscribe`, JSON.stringify(data));
  } catch (err) {
    throw err;
  }

  this.emit('rawSendData', data);

  // 封装等待响应的Promise
  return new Promise((resolve, reject) => {
    this.waitingMessageIds[messageId] = {
      resolve,
      reject,
      timeout: setTimeout(() => {
        delete this.waitingMessageIds[messageId];
        reject(new Error('Timeout'));
      }, 20000)
    };
  });
}

注:你当前用数组存waitingMessageIds是不合理的,字符串类型的messageId作为键应该用普通对象{}或Map存储,避免出现异常问题。


2. 改造message事件回调

将原调用callback的逻辑替换为调用Promise的resolve方法:

this.client.on('message', (topic, message) => {
  if (!message) return;
  try {
    message = JSON.parse(message.toString());
  } catch (err) {
    this.emit('error', `JSON parse error: ${err}`);
    return;
  }
  if (message.header.from && !message.header.from.includes(this.uuid)) return;

  const msgId = message.header.messageId;
  if (this.waitingMessageIds[msgId]) {
    clearTimeout(this.waitingMessageIds[msgId].timeout);
    // 把响应结果传给Promise的resolve
    this.waitingMessageIds[msgId].resolve(message.payload || message);
    delete this.waitingMessageIds[msgId];
  } else if (message.header.method === 'PUSH') {
    const namespace = message.header ? message.header.namespace : '';
    this.log('Found message');
    this.emit('data', namespace, message.payload || message);
  }
  this.emit('rawData', message);
});

3. 改造controlToggleX方法

改为async方法直接返回publishMessage的Promise结果:

async controlToggleX(channel, onoff) {
  const payload = { togglex: { channel, onoff: onoff ? 1 : 0 } };
  return this.publishMessage('SET', 'Appliance.Control.ToggleX', payload);
}

4. 上层调用示例

使用async/await语法获取结果:

// 调用逻辑需要放在async函数内部
try {
  const res = await super.controlToggleX(channel, value);
  this.log(`Toggle Response: res: ${JSON.stringify(res)}`);
} catch (err) {
  this.log(`Toggle Response: err: ${err}`);
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.23 22:57:02