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

在Rust中如何用阻塞函数实现单线程并发?

在无阻塞API的情况下实现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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 16:50:42