如何设计异步Rust的MAVLink库实现无消息丢失的发布/订阅?
问题描述
我正在开发一款用于MAVLink协议的异步Rust库,用于和无人机通信。通信流程为:消息发布任务从底层UDP连接接收字节流,解析为MAVLink消息后,分发给心跳任务、命令协议任务等业务处理任务。
目前我用tokio::sync::watch通道做消息发布,已经知道这不是合理选择,也欢迎针对这个问题的建议。
每个任务的循环逻辑大致如下:
loop { let message = wait_for(message_id).await; process_message(message) }
任务通过wait_for()从发布者接收消息,接着处理消息。现在的核心问题是:任务执行process_message时没法继续等待新消息,导致部分消息丢失。
我自己考虑过两种解决方案:
- 用注册回调的方式解决,但我希望给终端用户提供基于Future的API,或许可以在内部用回调?
- 给每个任务创建环形缓冲区,缓存一定数量的消息(比如最多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) { // 处理心跳消息逻辑 }
三、避免消息丢失的核心原则
- 分离接收与处理:绝对不要在消息接收的主线任务中同步执行耗时的
process_message,必须把处理逻辑放到独立的异步任务中,让接收流程持续运行 - 合理设置缓冲区:根据无人机消息的发送频率设置足够大的通道缓冲区,避免生产速度过快导致消息溢出
- 协议层确认(可选):对于关键消息(如命令确认),可在MAVLink协议层面实现消息确认机制,确保无人机收到并处理消息,这是从协议层面补充避免丢包的手段
四、参考设计模式与资源
- 发布-订阅模式:这是当前场景的核心模式,tokio的broadcast通道就是该模式的成熟实现
- Actor模型:如果后续任务逻辑复杂,可考虑用actix或tokio-actor等库,每个任务作为独立Actor处理消息,框架会自动维护消息队列
- 现有Rust MAVLink库:可以参考mavlink-rs等成熟库的异步实现思路,借鉴其消息分发与任务调度的设计
内容的提问来源于stack exchange,提问作者brkydnc
相关产品推荐
相关产品推荐

