如何用std::sync::mpsc::channel实现异步任务向同步环境传数据
Rust同步-异步交互中mpsc Sender的编译错误解决
问题描述
程序以同步逻辑为主,仅部分操作需调用私有库的异步接口。创建Tokio Runtime并spawn异步任务,尝试通过std::sync::mpsc的阻塞recv()将数据从异步任务传递到同步环境,但代码无法编译。
问题代码
use std::sync::mpsc::{channel, Sender}; struct Worker { chan: Sender<bool> } impl Worker { async fn do_work(&self) { loop { // 执行其他异步操作 self.chan.send(true).unwrap(); } } } fn main() { let (tx, rx) = channel::<bool>(); let rt = tokio::runtime::Builder::new_current_thread() .enable_all() .build() .unwrap(); let _ = rt.spawn(async move { let worker = Worker { chan: tx }; worker.do_work().await; }); println!("received {}", rx.recv().unwrap()); }
编译错误
within `Worker`, the trait `Sync` is not implemented for `std::sync::mpsc::Sender<bool>` note: future is not `Send` as this value is used across an await
约束条件
- 仅向单个线程传递channel,克隆Sender无法解决问题
- 接收端为同步环境,无法使用异步channel
- Worker对象需执行多轮
send()及其他异步操作,且需以移动方式传递给其他库函数
错误原因
std::sync::mpsc::Sender仅实现Send trait(可跨线程移动),但未实现Sync trait(不可跨线程共享引用)。do_work方法使用&self参数,且在await点后仍持有该引用,Tokio要求被spawn的future必须是Send,这要求&Worker是Send——而&T是Send的前提是T实现Sync,因此Worker因包含非Sync的Sender而无法满足编译要求。
解决方案
方案1:用Arc<Mutex<Sender<T>>>封装Sender
将Sender包裹在Arc<Mutex>中,让Worker满足Sync要求。由于整个流程仅单个线程操作Sender,Mutex不会产生实际性能开销,同时保留Worker被多次调用、移动传递的能力:
use std::sync::{mpsc::{channel, Sender}, Arc, Mutex}; struct Worker { chan: Arc<Mutex<Sender<bool>>> } impl Worker { async fn do_work(&self) { loop { // 执行其他异步操作 self.chan.lock().unwrap().send(true).unwrap(); } } } fn main() { let (tx, rx) = channel::<bool>(); let rt = tokio::runtime::Builder::new_current_thread() .enable_all() .build() .unwrap(); let _ = rt.spawn(async move { let worker = Worker { chan: Arc::new(Mutex::new(tx)) }; // 支持将worker移动传递给外部库函数 some_external_function(worker).await; }); println!("received {}", rx.recv().unwrap()); } // 模拟外部库函数,接收移动的Worker async fn some_external_function(worker: Worker) { worker.do_work().await; }
方案2:让do_work获取Worker所有权
如果Worker仅需被调用一次do_work,可将方法签名改为接收self所有权,这样无需满足Sync要求:
use std::sync::mpsc::{channel, Sender}; struct Worker { chan: Sender<bool> } impl Worker { async fn do_work(self) { loop { // 执行其他异步操作 self.chan.send(true).unwrap(); } } } fn main() { let (tx, rx) = channel::<bool>(); let rt = tokio::runtime::Builder::new_current_thread() .enable_all() .build() .unwrap(); let _ = rt.spawn(async move { let worker = Worker { chan: tx }; worker.do_work().await; }); println!("received {}", rx.recv().unwrap()); }
内容的提问来源于stack exchange,提问作者Boris Mulder
相关产品推荐
相关产品推荐

