不使用tokio-mutex和channel,如何解决Future跨线程安全发送问题?
问题核心分析
你的代码编译报错的本质原因有两个:
std::sync::Mutex的锁守卫(MutexGuard)未实现Sendtrait:Tokio的任务调度器可能在await点将future转移到不同线程执行,持有非Send的锁守卫会违反线程安全要求,导致"future cannot be sent between threads safely"错误。- 生命周期不匹配:函数
two中,异步写操作的future持有了已被销毁的锁守卫引用,导致"borrowed value does not live long enough"错误。
除了你提到的两种方案,还有以下几种可行的解决方式:
1. 使用async-lock库的异步互斥锁
async-lock是轻量级异步同步原语库,它的Mutex与Tokio运行时完全兼容,锁守卫实现了Send,可以安全地在await期间持有。
修改后代码示例:
use async_lock::Mutex; // 需要在Cargo.toml添加async-lock依赖 use std::sync::Arc; use tokio::net::TcpStream; #[tokio::test] async fn test_send_arc_mutex_tcps() { async fn one(tcps: Arc<Mutex<TcpStream>>, src: String) { let mut stream = tcps.lock().await; stream.write(src.as_bytes()).await.unwrap(); } async fn two(tcps: Arc<Mutex<TcpStream>>, src: String) { let mut stream = tcps.lock().await; let _ = stream.write(src.as_bytes()).await; } let tcps = TcpStream::connect("192.168.1.10:8080").await.unwrap(); let tcps = Arc::new(Mutex::new(tcps)); let datum = [ "hello".to_string(), "rustl".to_string(), "tokio".to_string(), ]; for data in datum.iter() { let tcps1 = tcps.clone(); let tcps2 = tcps.clone(); let data1 = data.clone(); let data2 = data.clone(); tokio::spawn(async move { one(tcps1, data1).await; two(tcps2, data2).await; }); } }
2. 使用tokio::sync::Semaphore实现互斥访问
如果仅需保证同一时间只有一个任务操作TcpStream,可用信号量(Semaphore)替代互斥锁。初始化许可数为1的Semaphore,任务获取许可后再操作流,操作完成后自动释放许可:
use std::sync::Arc; use tokio::net::TcpStream; use tokio::sync::Semaphore; #[tokio::test] async fn test_send_arc_semaphore_tcps() { async fn one(tcps: Arc<TcpStream>, sem: Arc<Semaphore>, src: String) { let _permit = sem.acquire().await.unwrap(); tcps.try_clone().unwrap().write(src.as_bytes()).await.unwrap(); } async fn two(tcps: Arc<TcpStream>, sem: Arc<Semaphore>, src: String) { let _permit = sem.acquire().await.unwrap(); let _ = tcps.try_clone().unwrap().write(src.as_bytes()).await; } let tcps = TcpStream::connect("192.168.1.10:8080").await.unwrap(); let tcps = Arc::new(tcps); let sem = Arc::new(Semaphore::new(1)); let datum = [ "hello".to_string(), "rustl".to_string(), "tokio".to_string(), ]; for data in datum.iter() { let tcps1 = tcps.clone(); let tcps2 = tcps.clone(); let sem1 = sem.clone(); let sem2 = sem.clone(); let data1 = data.clone(); let data2 = data.clone(); tokio::spawn(async move { one(tcps1, sem1, data1).await; two(tcps2, sem2, data2).await; }); } }
这里调用try_clone()获取TcpStream的克隆(Tokio的TcpStream内部共享文件描述符,支持安全克隆),每个任务持有独立的流实例,通过信号量保证写操作的互斥性。
3. 封装同步写操作(应急方案,不推荐)
如果一定要使用std::sync::Mutex,可以将异步写操作封装到同步闭包里,用tokio::task::spawn_blocking在阻塞线程中执行。但这种方式会浪费线程资源,异步写的非阻塞优势被抵消,仅适合低并发场景:
use std::sync::Arc; use std::sync::Mutex; use tokio::net::TcpStream; use tokio::task; #[tokio::test] async fn test_send_arc_std_mutex_tcps() { async fn one(tcps: Arc<Mutex<TcpStream>>, src: String) { task::spawn_blocking(move || { let mut stream = tcps.lock().unwrap(); // 注意:此处使用同步write_all,阻塞线程中无法使用异步await stream.write_all(src.as_bytes()).unwrap(); }).await.unwrap(); } async fn two(tcps: Arc<Mutex<TcpStream>>, src: String) { let _ = task::spawn_blocking(move || { let mut stream = tcps.lock().unwrap(); stream.write_all(src.as_bytes()) }).await.unwrap(); } let tcps = TcpStream::connect("192.168.1.10:8080").await.unwrap(); let tcps = Arc::new(Mutex::new(tcps)); let datum = [ "hello".to_string(), "rustl".to_string(), "tokio".to_string(), ]; for data in datum.iter() { let tcps1 = tcps.clone(); let tcps2 = tcps.clone(); let data1 = data.clone(); let data2 = data.clone(); tokio::spawn(async move { one(tcps1, data1).await; two(tcps2, data2).await; }); } }
内容的提问来源于stack exchange,提问作者progquester
相关产品推荐
相关产品推荐

