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

Node.js接入GCP IoT Core MQTT桥无法同时发布订阅消息问题

GCP Cloud IoT MQTT桥同时发布订阅问题修复

当然有大量开发者在Node.js里实现过同时收发的功能,你这问题不是平台不支持,是拼官方两段示例代码的时候漏了关键逻辑。

你代码里的核心问题

  • 订阅调用时机完全错误:你把client.subscribe()写在了connect事件监听的外面,刚执行完mqtt.connect()的时候TLS连接还没握手完成,这时候发的订阅请求会直接被客户端丢弃,根本不会传到GCP服务端,订阅自然不生效。
  • JWT刷新重连后没有重新订阅:Cloud IoT Core不支持MQTT持久会话,你重连时设置clean: true是对的,但新客户端连上之后完全没重新执行订阅操作,之前的订阅关系全失效,肯定收不到下行的config和command消息。
  • 重复绑定事件导致逻辑冲突:你在publishAsync的重连逻辑里又给新客户端重复绑了一套connect/message/error事件,跑久了不仅内存泄漏,新写的message回调还会和外层的消息处理逻辑抢消息,导致处理异常。
  • 示例自带的退出逻辑不适合长连接场景:官方发布示例是一次性发N条消息就主动断连退出的demo,发够numMessages就会执行client.end()断连,连都断了当然收不到订阅消息。
  • 变量混用容易抛错:代码里一会用全局定义的deviceId/registryId,一会从argv里取参数,作用域不对的时候会直接抛错打断整个流程。

可直接用的修复方案

  1. 先把公共逻辑抽出来,不要重复写代码,订阅必须放在连接成功的回调里执行,首次连接、重连都走同一套初始化逻辑:
// 先初始化全局状态
let publishChainInProgress = false;
const MINIMUM_BACKOFF_TIME = 1;
const MAXIMUM_BACKOFF_TIME = 32;
let shouldBackoff = false;
let backoffTime = MINIMUM_BACKOFF_TIME;
let client;
const iatTime = parseInt(Date.now() / 1000);
const mqttTopic = `/devices/${deviceId}/${messageType}`;

// 抽离连接成功后的统一初始化:先订阅,再启动发布
const initAfterConnect = (currentClient) => {
  currentClient.subscribe(`/devices/${deviceId}/config`, {qos: 1});
  currentClient.subscribe(`/devices/${deviceId}/commands/#`, {qos: 0});
  if (!publishChainInProgress) {
    publishAsync(mqttTopic, currentClient, iatTime, 1, numMessages, connectionArgs);
  }
};

// 抽离统一的事件绑定,避免重复绑定
const bindClientEvents = (currentClient) => {
  currentClient.on('connect', success => {
    console.log('mqtt connected');
    if (!success) {
      console.log('connect failed');
      return;
    }
    initAfterConnect(currentClient);
  });

  currentClient.on('close', () => {
    console.log('mqtt connection closed');
    shouldBackoff = true;
  });

  currentClient.on('error', err => {
    console.log('mqtt error', err);
  });

  currentClient.on('message', (topic, message) => {
    let messageStr = 'Message received: ';
    if (topic === `/devices/${deviceId}/config`) {
      messageStr = 'Config message received: ';
    } else if (topic.startsWith(`/devices/${deviceId}/commands`)) {
      messageStr = 'Command message received: ';
    }
    messageStr += Buffer.from(message, 'base64').toString('ascii');
    console.log(messageStr);
  });

  currentClient.on('packetsend', () => {});
};

// 首次初始化连接
client = mqtt.connect(connectionArgs);
bindClientEvents(client);
  1. 修改publishAsync逻辑,删掉重复的事件绑定,重连时复用公共逻辑,长期运行的话删掉发够消息就断连的退出逻辑:
const publishAsync = (
  mqttTopic,
  currentClient,
  iatTime,
  messagesSent,
  numMessages,
  connectionArgs
) => {
  // 长期同时收发场景下,注释掉原来发够消息就退出的逻辑
  // if (messagesSent > numMessages || backoffTime >= MAXIMUM_BACKOFF_TIME) {
  //   if (backoffTime >= MAXIMUM_BACKOFF_TIME) {
  //     console.log('Backoff time is too high. Giving up.');
  //   }
  //   console.log('Closing connection to MQTT. Goodbye!');
  //   currentClient.end();
  //   publishChainInProgress = false;
  //   return;
  // }

  publishChainInProgress = true;
  let publishDelayMs = 0;
  if (shouldBackoff) {
    publishDelayMs = 1000 * (backoffTime + Math.random());
    backoffTime *= 2;
    console.log(`Backing off for ${publishDelayMs}ms before publishing.`);
  }

  setTimeout(() => {
    // 统一用全局配置变量,不要混用argv避免作用域错误
    const payload = `${registryId}/${deviceId}-payload-${messagesSent}`;
    console.log('Publishing message:', payload);
    currentClient.publish(mqttTopic, payload, {qos: 1}, err => {
      if (!err) {
        shouldBackoff = false;
        backoffTime = MINIMUM_BACKOFF_TIME;
      }
    });

    const schedulePublishDelayMs = messageType === 'events' ? 1000 : 2000;
    setTimeout(() => {
      const secsFromIssue = parseInt(Date.now() / 1000) - iatTime;
      // JWT过期刷新逻辑
      if (secsFromIssue > tokenExpMins * 60) {
        iatTime = parseInt(Date.now() / 1000);
        console.log(`Refreshing JWT after ${secsFromIssue} seconds.`);
        currentClient.end(true);
        // 生成新的JWT
        connectionArgs.password = createJwt(projectId, privateKeyFile, algorithm);
        connectionArgs.clean = true;
        // 创建新连接,复用统一的事件绑定和初始化逻辑
        const newClient = mqtt.connect(connectionArgs);
        bindClientEvents(newClient);
        client = newClient;
        currentClient = newClient;
      }
      publishAsync(
        mqttTopic,
        currentClient,
        iatTime,
        messagesSent + 1,
        numMessages,
        connectionArgs
      );
    }, schedulePublishDelayMs);
  }, publishDelayMs);
};

改完之后就能正常同时发布消息、接收下行推送了。官方给的两段代码本来是分别对应“仅发布”“仅订阅”两个独立场景的最小demo,直接拼接很容易漏掉连接时序、重连恢复的必要逻辑。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.30 02:36:19