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里取参数,作用域不对的时候会直接抛错打断整个流程。
可直接用的修复方案
- 先把公共逻辑抽出来,不要重复写代码,订阅必须放在连接成功的回调里执行,首次连接、重连都走同一套初始化逻辑:
// 先初始化全局状态 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);
- 修改
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
相关产品推荐
相关产品推荐

