如何为rumqttc的阻塞迭代设置超时?
解决方案
要实现类似mpsc::recv_timeout的超时通知检查,你可以结合标准库的std::time::timeout和rumqttc同步客户端的Connection::next()方法(该方法会阻塞等待下一条通知),具体实现如下:
修改后的代码
// send it via mqtt use rumqttc::{MqttOptions, Client, QoS, Notification}; use std::time::{Duration, timeout}; use std::error::Error; fn main() -> Result<(), Box<dyn Error>> { let mut mqttoptions = MqttOptions::new("rumqtt-sync", "mqtt.example.com", 1883); mqttoptions.set_keep_alive(Duration::from_secs(5)); let (mut client, mut connection) = Client::new(mqttoptions, 10); // 假设json是已定义的有效负载 let json = b"{\"key\": \"value\"}"; let _ = client.publish("foo/bar", QoS::AtLeastOnce, false, json); // 3秒内检查通知,超时则继续执行后续逻辑 match timeout(Duration::from_secs(3), connection.next()) { // 成功获取到通知 Ok(Ok(notification)) => { println!("Notification = {:?}", notification); // 在这里添加通知处理逻辑 } // 获取通知时发生错误 Ok(Err(e)) => { eprintln!("Failed to receive notification: {:?}", e); } // 超时未收到通知,继续执行其他任务 Err(_) => { println!("No notifications received within 3 seconds, proceeding..."); // 在这里添加你需要执行的后续代码 } } Ok(()) }
关键逻辑说明
connection.next():同步阻塞等待下一条MQTT通知,直到有通知到达、连接出错或断开;timeout(Duration::from_secs(3), ...):将阻塞操作包裹,若超过3秒仍无结果,会返回Err(std::time::Elapsed),此时即可跳过等待,继续执行后续代码;- 三个分支分别处理不同场景,你可以根据需求在对应分支添加业务逻辑。
持续循环检查的场景
如果需要反复检查通知(每次等待3秒,超时后继续下一轮检查),可以把逻辑放到循环中:
loop { match timeout(Duration::from_secs(3), connection.next()) { Ok(Ok(notification)) => { println!("Notification = {:?}", notification); // 可根据通知类型决定是否退出循环,比如收到发布ACK后终止检查 } Ok(Err(e)) => { eprintln!("Connection error: {:?}", e); break; } Err(_) => { println!("Timeout, checking for notifications again..."); // 执行周期性任务 } } }
内容的提问来源于stack exchange,提问作者fadedbee
相关产品推荐
相关产品推荐

