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

如何在Tokio异步任务中正确使用DuplexStream实现并发读写?

解决Tokio DuplexStream并发读写的锁阻塞问题

问题核心原因

Tokio的DuplexStream本身已经实现了异步安全的并发读写支持,内部通过同步原语处理了读写之间的竞争问题。你额外用Arc<Mutex<DuplexStream>>包裹完全是画蛇添足——读任务在执行read().await期间会长期持有锁,直接导致写任务无法获取锁执行写入操作。

正确实现方式

直接通过Arc共享DuplexStream即可,不需要额外的Mutex包裹。因为DuplexStream满足Send + Sync trait(只要初始化时指定的缓冲区大小是固定值),可以安全地在多个异步任务中并发调用读写方法。

示例代码

use tokio::io::{duplex, AsyncReadExt, AsyncWriteExt};
use std::sync::Arc;

#[tokio::main]
async fn main() {
    // 创建双向流,缓冲区大小64KB
    let (upstream, dwstream) = duplex(64 * 1024);
    let shared_stream = Arc::new(dwstream);
    let shared_stream_clone = shared_stream.clone();

    // 异步写任务
    tokio::spawn(async move {
        let message = b"Hello from write task!";
        // 直接调用write_all,无需锁
        shared_stream.write_all(message).await.unwrap();
        println!("Write task completed: sent {}", String::from_utf8_lossy(message));
    });

    // 异步读任务
    let read_handle = tokio::spawn(async move {
        let mut buffer = [0; 256];
        // 直接调用read,无需锁
        let bytes_read = shared_stream_clone.read(&mut buffer).await.unwrap();
        println!("Read task completed: received {}", String::from_utf8_lossy(&buffer[..bytes_read]));
    });

    // 等待读任务完成
    read_handle.await.unwrap();
}

额外说明

如果你的业务逻辑需要对DuplexStream的操作做更复杂的同步控制(比如读写操作之间有业务层面的依赖),可以使用tokio::sync::Mutex(Tokio提供的异步锁,而非标准库的std::sync::Mutex),但单纯的并发读写场景下完全不需要。

内容的提问来源于stack exchange,提问作者progquester

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.19 06:00:17