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

如何设计异步Rust的MAVLink库实现无消息丢失的发布/订阅?

问题描述

我正在开发一款用于MAVLink协议的异步Rust库,用于和无人机通信。通信流程为:消息发布任务从底层UDP连接接收字节流,解析为MAVLink消息后,分发给心跳任务、命令协议任务等业务处理任务。

目前我用tokio::sync::watch通道做消息发布,已经知道这不是合理选择,也欢迎针对这个问题的建议。

每个任务的循环逻辑大致如下:

loop {
    let message = wait_for(message_id).await;
    process_message(message)
}

任务通过wait_for()从发布者接收消息,接着处理消息。现在的核心问题是:任务执行process_message时没法继续等待新消息,导致部分消息丢失。

我自己考虑过两种解决方案:

  1. 用注册回调的方式解决,但我希望给终端用户提供基于Future的API,或许可以在内部用回调?
  2. 给每个任务创建环形缓冲区,缓存一定数量的消息(比如最多32条),允许少量丢包,但这会不会太复杂?

另外,当前设计可能存在缺陷,也欢迎新的设计思路。请问该如何设计这个库,让任务永远或几乎不丢失消息?有没有对应的设计模式或参考资源可以借鉴?


解决方案与设计思路

一、替换通道类型:用tokio::sync::broadcast替代watch

tokio::sync::watch本质是单生产者单消费者且仅保留最新值的通道,完全不适合多任务订阅、需要保留历史消息的场景。换成broadcast是最直接的改进:

  • 支持单生产者多消费者,每个消费者独立维护自己的消息游标
  • 每个任务可创建专属接收器,在处理消息的间隙异步等待新消息,不会阻塞消息接收流程
  • 可设置通道容量,当消息生产速度超过消费速度时,旧消息会被覆盖,避免内存溢出

示例改造:

// 发布端初始化
let (tx, _) = tokio::sync::broadcast::channel(64); // 根据消息频率设置合适的缓冲区大小

// 每个任务的接收器
let mut rx = tx.subscribe();

loop {
    // 异步等待新消息,无需阻塞处理流程
    match rx.recv().await {
        Ok(message) => {
            // 将消息处理逻辑放到独立异步任务中,不阻塞接收
            tokio::spawn(process_message(message));
        }
        Err(e) => {
            // 处理订阅错误,比如生产者关闭
            eprintln!("订阅出错: {}", e);
            break;
        }
    }
}

二、基于Future的API设计:用Stream封装消息流

如果要给用户提供更友好的Future/Stream风格API,可以把每个消息类型的订阅封装成Stream:

  • 利用tokio_stream::wrappers::BroadcastStream将broadcast接收器转换为Stream
  • 提供subscribe_message(message_id)这类方法,返回过滤后的目标消息Stream
  • 用户可通过while let Some(msg) = stream.next().await消费消息,同时用tokio::spawn并行处理消息,避免阻塞接收

示例:

use tokio_stream::{Stream, StreamExt};
use tokio::sync::broadcast;

// 定义MAVLink消息类型
#[derive(Debug)]
pub struct MavlinkMessage {
    pub id: u32,
    // 其他消息字段
}

pub struct MavlinkClient {
    tx: broadcast::Sender<MavlinkMessage>,
}

impl MavlinkClient {
    pub fn subscribe_message(&self, message_id: u32) -> impl Stream<Item = MavlinkMessage> {
        let mut rx = self.tx.subscribe();
        tokio_stream::wrappers::BroadcastStream::new(rx)
            .filter(move |msg| msg.as_ref().map(|m| m.id == message_id).unwrap_or(false))
            .map(|msg| msg.unwrap())
    }
}

// 用户侧使用示例
async fn user_code(client: MavlinkClient) {
    const HEARTBEAT_ID: u32 = 0; // MAVLink心跳消息ID
    let mut heartbeat_stream = client.subscribe_message(HEARTBEAT_ID);
    tokio::spawn(async move {
        while let Some(heartbeat) = heartbeat_stream.next().await {
            // 异步处理心跳,不阻塞Stream接收
            tokio::spawn(process_heartbeat(heartbeat));
        }
    });
}

async fn process_heartbeat(msg: MavlinkMessage) {
    // 处理心跳消息逻辑
}

三、避免消息丢失的核心原则

  1. 分离接收与处理:绝对不要在消息接收的主线任务中同步执行耗时的process_message,必须把处理逻辑放到独立的异步任务中,让接收流程持续运行
  2. 合理设置缓冲区:根据无人机消息的发送频率设置足够大的通道缓冲区,避免生产速度过快导致消息溢出
  3. 协议层确认(可选):对于关键消息(如命令确认),可在MAVLink协议层面实现消息确认机制,确保无人机收到并处理消息,这是从协议层面补充避免丢包的手段

四、参考设计模式与资源

  • 发布-订阅模式:这是当前场景的核心模式,tokio的broadcast通道就是该模式的成熟实现
  • Actor模型:如果后续任务逻辑复杂,可考虑用actix或tokio-actor等库,每个任务作为独立Actor处理消息,框架会自动维护消息队列
  • 现有Rust MAVLink库:可以参考mavlink-rs等成熟库的异步实现思路,借鉴其消息分发与任务调度的设计

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 00:52:04