Tokio mpsc Receiver非阻塞接收:单线程处理非线程安全外设收发方案
问题根因
device_query库内部持有X11原生裸指针,未实现Send trait。而tokio默认多线程调度器会在工作线程间转移异步任务,要求所有跨await存活的变量必须满足Send约束,因此编译报错。
解决方案
方案1:使用tokio本地任务(保留原有异步逻辑)
将设备任务绑定到单线程执行,不需要Send约束,可正常使用await:
- 将
tokio::task::spawn替换为tokio::task::spawn_local,该API创建的任务只会在生成它的线程上运行,不会跨线程转移。 - 调整tokio运行时为单线程模式,确保本地任务可正常调度。
修改后代码:
use device_query::{DeviceState, Keycode}; use std::time::Duration; use tokio; use tokio::sync::mpsc::{Receiver, Sender}; use tokio::time::timeout; // 指定tokio使用单线程运行时 #[tokio::main(flavor = "current_thread")] async fn main() { // 设备循环产出的按键事件通道 let (key_tx, mut key_rx) = tokio::sync::mpsc::channel(32); // 发往设备循环的命令通道 let (dev_tx, mut dev_rx) = tokio::sync::mpsc::channel(32); start_device_loop(60, key_tx, dev_rx); println!("Waiting for key presses"); while let Some(k) = key_rx.recv().await { match k { Some(ch) => match ch { Keycode::Q => dev_tx.clone().try_send(String::from("Quit!")).expect("Could not send command"), ch => println!("{}", ch), }, _ => (), } } println!("Done.") } /// 启动tokio本地任务,轮询指定设备并通过传入的mpsc sender发送按键事件 pub fn start_device_loop(hz: u32, tx: Sender<Option<Keycode>>, mut rx: Receiver<String>) { let poll_wait = 1000 / hz; let poll_wait = Duration::from_millis(poll_wait as u64); // spawn_local创建的任务不需要实现Send tokio::task::spawn_local(async move { let dev = DeviceState::new(); loop { let mut keys = dev.query_keymap(); match keys.len() { 0 => (), 1 => tx.clone().try_send(Some(keys.remove(0))).unwrap(), _ => println!("So many keys..."), } match timeout(poll_wait, rx.recv()).await { Ok(Some(cmd)) => println!("Command '{}' received.", cmd), // 命令通道关闭时退出任务 Ok(None) => break, Err(_) => (), }; } }); }
方案2:使用原生std线程(完全不使用async/await)
如果不需要异步能力,直接开启系统原生线程处理设备逻辑,用非阻塞方法读取命令:
pub fn start_device_loop(hz: u32, tx: Sender<Option<Keycode>>, mut rx: Receiver<String>) { let poll_wait = 1000 / hz; let poll_wait = Duration::from_millis(poll_wait as u64); std::thread::spawn(move || { let dev = DeviceState::new(); loop { let mut keys = dev.query_keymap(); match keys.len() { 0 => (), 1 => tx.clone().try_send(Some(keys.remove(0))).unwrap(), _ => println!("So many keys..."), } // 非阻塞读取命令,不会阻塞轮询 match rx.try_recv() { Ok(cmd) => println!("Command '{}' received.", cmd), Err(tokio::sync::mpsc::error::TryRecvError::Empty) => (), // 命令通道关闭时退出线程 Err(tokio::sync::mpsc::error::TryRecvError::Disconnected) => break, }; // 睡眠固定时长控制轮询频率 std::thread::sleep(poll_wait); } }); }
该方案main函数无需修改,完全不涉及异步语法,也满足单线程操作设备的要求。
内容的提问来源于stack exchange,提问作者superstator
相关产品推荐
相关产品推荐

