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

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);
});

问题分析

  1. 锁持有时间过长:读线程调用read_packet时,read_exact是阻塞IO操作,会长期持有Mutex锁;写线程每隔1秒尝试获取锁,若读线程一直阻塞等待数据,写线程将永远拿不到锁,形成死锁。
  2. 阻塞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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 00:24:59