Tokio任务线程共享疑问:阻塞任务为何抢占核心线程?
问题分析:Tokio异步任务阻塞引发的调度异常
问题代码
use std::time::Duration; use tokio::sync::mpsc; async fn send(tx: mpsc::Sender<u32>) { loop { std::thread::sleep(Duration::from_secs(1)); tx.try_send(42).unwrap(); println!("sent"); } } async fn recv(mut rx: mpsc::Receiver<u32>) { while let Some(_) = rx.recv().await { println!("recv"); } } #[tokio::main] async fn main() { let (tx, rx) = mpsc::channel(4); tokio::spawn(send(tx.clone())); tokio::spawn(recv(rx)); // let it run for a bit to see the channel fill up and send task to panic tokio::time::sleep(Duration::from_secs(100)).await; }
疑问点
这段代码故意在异步任务中使用std::thread::sleep(无.await点),运行后发现接收任务直到通道填满导致发送任务panic时才执行recv().await。明明电脑是多核,且Tokio运行时默认核心线程数等于CPU核心数,为何两个任务没分配到不同核心线程?是否所有任务先共享一个核心线程,直到运行时判定需要新线程才分配?原本以为任务1占用核心线程1,任务2占用核心线程2(已知示例违反异步代码不阻塞的基本要求)。
原因解析
- Tokio的工作窃取调度器默认不会主动为每个新任务分配独立核心线程。调度器会先把任务放到某个核心线程的本地队列中,只有当该线程的本地队列任务处理不过来,或者出现阻塞时,才会触发工作窃取,或者考虑创建新线程(针对阻塞情况)。
std::thread::sleep是阻塞整个线程的同步睡眠,而非Tokio提供的tokio::time::sleep(异步睡眠会在等待时主动让出线程)。当send任务在核心线程上执行std::thread::sleep时,这个线程被完全占死,无法切换到队列里的其他任务。- 初始调度时,Tokio可能把两个任务都放到了同一个核心线程的队列中,
send任务一直阻塞着该线程,导致recv任务根本没机会被调度执行,直到send任务因通道满panic,线程才得以释放,recv任务才开始运行。 - 哪怕是多核环境,Tokio不会一开始就把任务分散到所有核心线程,而是基于负载情况动态调度。只有当现有线程无法处理任务(比如出现阻塞、任务队列积压),才会利用其他核心线程或者创建新线程。
修正方案
把std::thread::sleep替换为Tokio的异步睡眠tokio::time::sleep(Duration::from_secs(1)).await,这样任务在等待时会主动让出线程,调度器就能切换到recv任务执行,两者可以正常并发:
use std::time::Duration; use tokio::sync::mpsc; use tokio::time; async fn send(tx: mpsc::Sender<u32>) { loop { time::sleep(Duration::from_secs(1)).await; tx.try_send(42).unwrap(); println!("sent"); } } async fn recv(mut rx: mpsc::Receiver<u32>) { while let Some(_) = rx.recv().await { println!("recv"); } } #[tokio::main] async fn main() { let (tx, rx) = mpsc::channel(4); tokio::spawn(send(tx.clone())); tokio::spawn(recv(rx)); tokio::time::sleep(Duration::from_secs(100)).await; }
内容的提问来源于stack exchange,提问作者jsstuball
相关产品推荐
相关产品推荐

