如何在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
相关产品推荐
相关产品推荐

