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

如何为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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.12 14:22:44