Rust并发读写阻塞Socket时Arc引发死锁的解决方案咨询
Rust阻塞Socket并发读写的死锁问题与解决方案
问题背景
在Rust中使用阻塞Socket实现线程分离的并发读写时,若通过Arc<Mutex>共享Transport或Socket实例,极易出现死锁:读线程因阻塞IO长期持有锁,导致写线程无法获取锁执行写操作,反之亦然。
示例代码
Rust客户端程序
use std::io::{Error, Read, Write}; use std::net::{TcpStream, ToSocketAddrs}; use std::sync::{Arc, Mutex}; use std::time::Duration; use tokio::time::sleep; #[derive(Debug)] pub struct TcpTransport { pub conn: TcpStream, } impl TcpTransport { fn send_packet(&mut self, data: &[u8]) -> Result<(), Error> { self.conn.write_all(data)?; Ok(()) } fn read_packet(&mut self) -> Result<Vec<u8>, Error> { let mut lenbuf = [0u8; 2]; self.conn.read_exact(&mut lenbuf)?; let length = (lenbuf[0] as usize) << 8 | lenbuf[1] as usize; let mut databuf = vec![0u8; length]; self.conn.read_exact(&mut databuf)?; Ok(databuf) } fn close(&self) -> Result<(), Error> { println!("close transport"); self.conn.shutdown(std::net::Shutdown::Both)?; Ok(()) } } pub fn new_tcp_transport(host: &str, port: u16) -> Result<TcpTransport, Error> { let socket_addr = (host, port).to_socket_addrs()?.next().unwrap(); let conn = TcpStream::connect_timeout(&socket_addr, Duration::from_secs(10))?; Ok(TcpTransport { conn }) } #[tokio::main] async fn main() { let tsport = new_tcp_transport("127.0.0.1", 4000).unwrap(); let tsport = Arc::new(Mutex::new(tsport)); let tsport1 = tsport.clone(); let tsport2 = tsport.clone(); tokio::spawn(async move { println!("read worker started"); 'read_loop: loop { println!("read exec here==="); let packet = match tsport1.lock().unwrap().read_packet() { Ok(packet) => packet, Err(err) => { eprintln!("Transport read packet error: {:?}", err); break 'read_loop; } }; println!("read packet{:?}", packet); } println!("read worker stoped"); }); let _ = tokio::spawn(async move { println!("write worker started"); 'writeloop: loop { sleep(Duration::from_millis(1000)).await; println!("write exec here==="); let mut ts = tsport2.lock().unwrap(); let data = vec![0x1,0x2,0x3,0x4]; if let Err(er) = ts.send_packet(&data) { eprintln!("Transport write packet error: {:?}", er); break 'writeloop; } println!("send packet success."); } println!("write worker stoped"); }).await; }
Node.js服务端程序
var net = require('net'); net.createServer(function(socket){ socket.on('data', function(data){ console.log("server recv data:",data); }); }).listen(4000); console.log('server listen 127.0.0.1:4000');
Cargo.toml配置
[package] name = "readwrite" version = "0.1.0" edition = "2021" [dependencies] tokio = { version = "1.28.1", features = ["full"] }
Node.js客户端程序
var net = require('net'); var client = new net.Socket(); client.connect(4000, '127.0.0.1', function() { console.log('Connected'); setInterval(() => { var data = Buffer.from("hello server"); var datalen = data.length; console.log('client write==='); client.write(Buffer.concat([Buffer.from([datalen >> 8, datalen % 256]), data])); }, 1000); }); client.on('data', function(data) { console.log('client received: ',data); });
问题分析
- 锁持有时间过长:读线程调用
read_packet时,read_exact是阻塞IO操作,会长期持有Mutex锁;写线程每隔1秒尝试获取锁,若读线程一直阻塞等待数据,写线程将永远拿不到锁,形成死锁。 - 阻塞IO与异步Runtime冲突:在Tokio异步Runtime中运行阻塞IO,会占用工作线程,导致其他异步任务无法及时调度,加剧锁竞争问题。
解决方案
方案1:分离Socket读写半(阻塞场景推荐)
利用标准库TcpStream的split()方法,将Socket拆分为独立的读半(ReadHalf)和写半(WriteHalf),两者可分别通过Arc共享,各自使用独立的锁,彻底避免读写线程竞争同一锁资源。
修改后的核心代码示例:
#[tokio::main] async fn main() { let mut transport = new_tcp_transport("127.0.0.1", 4000).unwrap(); // 拆分读写半 let (read_half, write_half) = transport.conn.split(); let read_arc = Arc::new(Mutex::new(read_half)); let write_arc = Arc::new(Mutex::new(write_half)); // 读线程 let read_clone = read_arc.clone(); tokio::spawn(async move { println!("read worker started"); let mut buf_len = [0u8;2]; loop { match read_clone.lock().unwrap().read_exact(&mut buf_len) { Ok(_) => { let len = (buf_len[0] as usize) <<8 | buf_len[1] as usize; let mut data = vec![0u8; len]; if let Ok(_) = read_clone.lock().unwrap().read_exact(&mut data) { println!("read packet: {:?}", data); } else { break; } } Err(e) => { eprintln!("read error: {:?}", e); break; } } } println!("read worker stopped"); }); // 写线程 let write_clone = write_arc.clone(); let _ = tokio::spawn(async move { println!("write worker started"); loop { sleep(Duration::from_millis(1000)).await; // 补充协议要求的长度前缀 let data = vec![0x00, 0x04, 0x1,0x2,0x3,0x4]; if let Err(e) = write_clone.lock().unwrap().write_all(&data) { eprintln!("write error: {:?}", e); break; } println!("send packet success"); } println!("write worker stopped"); }).await; }
方案2:改用异步Socket(异步场景推荐)
使用Tokio提供的异步TcpStream,结合Tokio的异步互斥锁tokio::sync::Mutex,避免阻塞IO占用线程,同时异步锁会在等待时让出线程,不影响其他任务调度。
修改后的核心代码示例:
use tokio::net::TcpStream; use tokio::io::{AsyncReadExt, AsyncWriteExt}; use tokio::sync::Mutex; use std::sync::Arc; use std::time::Duration; #[derive(Debug)] pub struct AsyncTcpTransport { conn: TcpStream, } impl AsyncTcpTransport { async fn send_packet(&mut self, data: &[u8]) -> Result<(), std::io::Error> { self.conn.write_all(data).await?; Ok(()) } async fn read_packet(&mut self) -> Result<Vec<u8>, std::io::Error> { let mut lenbuf = [0u8; 2]; self.conn.read_exact(&mut lenbuf).await?; let length = (lenbuf[0] as usize) << 8 | lenbuf[1] as usize; let mut databuf = vec![0u8; length]; self.conn.read_exact(&mut databuf).await?; Ok(databuf) } } async fn new_async_tcp_transport(host: &str, port: u16) -> Result<AsyncTcpTransport, std::io::Error> { let addr = format!("{}:{}", host, port); let conn = TcpStream::connect(addr).await?; Ok(AsyncTcpTransport { conn }) } #[tokio::main] async fn main() { let transport = new_async_tcp_transport("127.0.0.1", 4000).await.unwrap(); let transport = Arc::new(Mutex::new(transport)); let read_clone = transport.clone(); tokio::spawn(async move { println!("read worker started"); loop { let mut transport = read_clone.lock().await; match transport.read_packet().await { Ok(packet) => println!("read packet: {:?}", packet), Err(e) => { eprintln!("read error: {:?}", e); break; } } } println!("read worker stopped"); }); let write_clone = transport.clone(); let _ = tokio::spawn(async move { println!("write worker started"); loop { tokio::time::sleep(Duration::from_millis(1000)).await; let mut transport = write_clone.lock().await; // 补充协议要求的长度前缀 let data = vec![0x00, 0x04, 0x1, 0x2, 0x3, 0x4]; if let Err(e) = transport.send_packet(&data).await { eprintln!("write error: {:?}", e); break; } println!("send packet success"); } println!("write worker stopped"); }).await; }
总结
- 若坚持使用阻塞Socket,优先拆分读写半,分离锁的作用域,避免读写线程竞争同一锁。
- 若使用异步Runtime,推荐改用异步Socket和异步锁,适配异步调度模型,从根源上避免阻塞导致的死锁问题。
内容的提问来源于stack exchange,提问作者Hai.Xu
相关产品推荐
相关产品推荐

