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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 06:15:28