在Rust中如何用阻塞函数实现单线程并发?
问题背景
我是Rust新手,想了解在没有非阻塞API的情况下,是否可以编写单线程并发应用。我基于UDP socket和tun设备(虚拟网络接口)编写了Linux环境下的VPN服务端代码,期望等待tun文件描述符的写入或UDP socket的数据包到达。代码如下:
loop { tokio::select! { _ = async { log::debug!("reading from interface"); tun.read(&mut buf1).unwrap(); } => { log::debug!("sending to client {}", &caddr); socket.send_to(&buf1, &caddr).unwrap(); }, _ = async { log::debug!("reading from client"); let (_amt, addr) = socket.recv_from(&mut buf2).unwrap(); caddr = addr.to_string(); } => { log::debug!("writting to interface"); tun.write(&mut buf2).unwrap(); } };
我原本认为tokio::select!能满足需求,其官方文档说明该宏:
等待多个并发分支,当第一个分支完成时返回,取消剩余分支
但运行后发现debug日志“reading from interface”从未执行,推测是因为阻塞的socket.recv_from()导致tokio::select!无法切换分支。虽然我找到了提供非阻塞API的tokio::net::UdpSocket和tun_tokio库可满足需求,但仍好奇在没有非阻塞API时,如何实现预期的单线程并发行为?
解决方案
在没有异步/非阻塞API的情况下,要让Tokio这类异步 runtime 实现单线程并发,核心是不能让阻塞操作占据Tokio的工作线程——因为Tokio的协作式调度依赖任务主动让出CPU。你的代码里tun.read()和socket.recv_from()都是阻塞调用,一旦其中一个执行,就会卡住整个线程,其他分支根本没机会运行。
下面是两种可行的实现思路:
1. 用spawn_blocking包装阻塞操作
Tokio提供的spawn_blocking可以把阻塞任务放到专门的阻塞线程池执行,避免阻塞异步工作线程。修改后的代码示例如下:
loop { tokio::select! { _ = tokio::task::spawn_blocking(move || { log::debug!("reading from interface"); tun.read(&mut buf1).unwrap(); }) => { log::debug!("sending to client {}", &caddr); socket.send_to(&buf1, &caddr).unwrap(); }, _ = tokio::task::spawn_blocking(move || { log::debug!("reading from client"); let (_amt, addr) = socket.recv_from(&mut buf2).unwrap(); caddr = addr.to_string(); }) => { log::debug!("writing to interface"); tun.write(&mut buf2).unwrap(); } }; }
注意:实际使用时需要处理好缓冲区的所有权问题,可能需要通过Arc<Mutex<T>>共享数据,或者调整变量生命周期以满足编译器要求。
2. 手动将FD设为非阻塞,配合Tokio事件监听
在Linux环境下,可以通过fcntl系统调用将socket和tun设备的文件描述符设置为非阻塞模式,再用Tokio的Interest监听FD的可读事件,手动处理非阻塞IO的返回结果。
步骤1:设置文件描述符为非阻塞
use libc::{fcntl, F_GETFL, F_SETFL, O_NONBLOCK}; // 将UDP socket设为非阻塞 let socket_fd = socket.as_raw_fd(); let socket_flags = unsafe { fcntl(socket_fd, F_GETFL) }; unsafe { fcntl(socket_fd, F_SETFL, socket_flags | O_NONBLOCK) }; // 将tun设备设为非阻塞 let tun_fd = tun.as_raw_fd(); let tun_flags = unsafe { fcntl(tun_fd, F_GETFL) }; unsafe { fcntl(tun_fd, F_SETFL, tun_flags | O_NONBLOCK) };
步骤2:监听事件并处理非阻塞IO
use tokio::io::{Interest, ready}; loop { tokio::select! { _ = ready(tun_fd, Interest::READABLE) => { log::debug!("reading from interface"); match tun.read(&mut buf1) { Ok(_) => { log::debug!("sending to client {}", &caddr); socket.send_to(&buf1, &caddr).unwrap(); }, Err(e) if e.kind() == std::io::ErrorKind::WouldBlock => continue, Err(e) => panic!("tun read error: {}", e), } }, _ = ready(socket_fd, Interest::READABLE) => { log::debug!("reading from client"); match socket.recv_from(&mut buf2) { Ok((_amt, addr)) => { caddr = addr.to_string(); log::debug!("writing to interface"); tun.write(&mut buf2).unwrap(); }, Err(e) if e.kind() == std::io::ErrorKind::WouldBlock => continue, Err(e) => panic!("socket recv error: {}", e), } } }; }
需要注意的是,这种方式需要手动处理非阻塞IO的WouldBlock错误,复杂度较高。如果有现成的异步IO库可用,优先选择官方或社区维护的异步版本会更简洁可靠。
内容的提问来源于stack exchange,提问作者shubmishra

