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

不使用tokio-mutex和channel,如何解决Future跨线程安全发送问题?

问题核心分析

你的代码编译报错的本质原因有两个:

  1. std::sync::Mutex的锁守卫(MutexGuard)未实现Send trait:Tokio的任务调度器可能在await点将future转移到不同线程执行,持有非Send的锁守卫会违反线程安全要求,导致"future cannot be sent between threads safely"错误。
  2. 生命周期不匹配:函数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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 23:15:43