MQTT v5 MessageExpiryInterval功能异常问题求助
MQTT v5 messageExpiryInterval 属性不生效问题
需求与问题
我需要实现带时间限制的消息发布,模拟设备进入隧道、无蜂窝信号等离线场景,让消息过期后Broker自动删除该消息,设备重新上线后不会接收到过期消息。根据MQTT v5协议,messageExpiryInterval属性应该能实现这一功能,但实际测试后不符合预期。
环境与配置
- 使用EMQX Broker,连接地址:
const connectUrl = 'mqtt://broker.emqx.io:1883'; - 发布者连接配置:
const client = mqtt.connect(connectUrl, { clientId: 'mqtt_dani_pub', protocolVersion: 5, keepalive: 1800, clean: true });
- 消息发布配置(设置30秒过期):
const options = { qos: 0, retain: false, properties: { messageExpiryInterval: 30 } };
- 订阅者使用相同协议版本与keepalive配置,仅clientId不同。
测试流程与异常结果
- 连接发布者与订阅者,确认双方均使用MQTT v5协议
- 订阅者在线时,发布
messageExpiryInterval=30的消息,订阅者正常接收 - 关闭订阅者Wi-Fi使其离线
- 发布
messageExpiryInterval=30的消息 - 等待120秒(远超30秒过期时间)
- 开启订阅者Wi-Fi重新上线,预期不会收到过期消息,但实际仍收到
协议相关疑问
查阅MQTT v5协议3.3.2.3.3章节,其中提到服务器发送给客户端的PUBLISH包,需要将Message Expiry Interval设为接收值减去消息在服务器等待的时间,怀疑这一逻辑未被正确执行是问题根源。
完整代码
发布者代码
import mqtt, { MqttClient } from 'mqtt'; import * as readline from 'node:readline' import { stdin, stdout } from 'process'; const connectUrl = 'mqtt://broker.emqx.io:1883'; const clientId = 'mqtt_dani_pub'; const topic = 'dani/test'; const subject = 'Publisher'; const rl = readline.createInterface({ input:stdin, output:stdout }); const client = mqtt.connect(connectUrl, { clientId, protocolVersion: 5, keepalive: 1800, clean: true }); client.on('connect', () => { console.log(`${subject} client connected..`) client.subscribe([topic], () => { console.log(`Subscribe to topic '${topic}'`); }) }); const options = { qos: 0, retain: false, properties: { messageExpiryInterval: 30 } }; const publishMsg = (message) => { client.publish(topic, `${clientId} - ${Date.now()} - ${message}`, options, (error) => { if (error) { console.error(error) } } ); }; const input2topic = () => { return new Promise(resolve => { rl.question(`send message to topic ${topic}: `, (input) => { if(input !== 'exit'){ console.log(`writing to topic ${topic}..`); publishMsg(input); resolve(true); } else{ console.log('exit...'); resolve(false); } }); }); } const main = async () => { publishMsg('first message'); let flag = true; while(flag){ await new Promise(resolve => setTimeout(resolve, 1000)); flag = await input2topic(); } rl.close(); client.end(); } main();
订阅者代码
import mqtt, { MqttClient } from 'mqtt'; const connectUrl = 'mqtt://broker.emqx.io:1883'; const clientId = 'mqtt_dani_sub'; const topic = 'dani/test'; const subject = 'Subscriber'; const client = mqtt.connect(connectUrl, { clientId, keepalive: 1800, protocolVersion: 5, }) client.on('connect', () => { console.log(`${subject} client connected`) client.subscribe([topic], {qos: 0}, () => { console.log(`Subscribe to topic '${topic}'`) }) }); client.on('message', (topic, payload, packet) => { console.log('\nReceived Message:', { ...packet, message: payload.toString(), msg_length: payload.toString().length, time: new Date(), }); });
内容的提问来源于stack exchange,提问作者DanielTal
相关产品推荐
相关产品推荐

