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

为何以异常方式调用Tokio时仅运行n-1个任务?

Rust Tokio并发任务异常问题分析

我最近在学习Rust并发/并行编程,编写不同并发实现的代码片段时碰到了Tokio的异常问题。我确认代码是误用了Tokio,程序没崩溃但运行不正常——总是少执行一个任务。我已经找到规避方法,但想搞懂内部原因,有两个疑问:

  1. 既然代码确实写错了,为什么还有部分任务能正常运行?
  2. 为什么偏偏总是少执行一个任务?

复现代码

#![allow(dead_code)]

// There are two ways to fix this code. They are marked with a "FIX" comment

use tokio;

const NUMBER_OF_THREADS: usize = 3;

fn main() {
    futures::executor::block_on(with_tokio_channels_fail_call_breaking(NUMBER_OF_THREADS));
    
    // FIX: use this line instead
    // with_tokio_channels_fail_call_fixing(NUMBER_OF_THREADS);
}

// Calling this function runs the code fine
fn with_tokio_channels_fail_call_fixing(number_of_threads: usize) {
    let rt = tokio::runtime::Runtime::new().unwrap();
    // When we use block_on on this fake-future (no .await inside), it runs
    // fine
    rt.block_on(tokio_channels_fail(number_of_threads));
}

// Calling this function runs (NUMBER_OF_THREADS - 1) threads fine,
// but the last one is never scheduled, causing the program to never exit
async fn with_tokio_channels_fail_call_breaking(number_of_threads: usize) {
    let rt = tokio::runtime::Runtime::new().unwrap();
    // Spawning and awaiting the fake-future has a really odd result!
    // I believe this is simply wrong usage of tokio. But WHY exactly?
    // I am interested in what goes wrong internally when we do this
    rt.spawn(tokio_channels_fail(number_of_threads)).await.unwrap();
}

// Note that this function is unnecessarily async
async fn tokio_channels_fail(number_of_threads: usize) {
    async fn cb(index: usize, num: usize, sender: std::sync::mpsc::Sender<(usize, usize)>) {
        let res = delay_thread_async_then_square(num, num).await;
        sender.send((index, res)).unwrap();
    }

    let (sender, receiver) = std::sync::mpsc::channel::<(usize, usize)>();
    let mut vec: Vec<usize> = (0..number_of_threads).into_iter().collect();

    let _x = (0..number_of_threads)
        .map(|num| {
            let sender = sender.clone();
            tokio::task::spawn(cb(num, num, sender))
        })
        .collect::<Vec<tokio::task::JoinHandle<()>>>();
    // need to drop the sender, because the iterator below will only complete once all senders are dropped
    drop(sender);

    // FIX: when uncommenting this line, the code also runs fine, but this is not a general solution.
    // But would it use the tokio runtime? Means: Would it be spread across multiple
    // kernel level threads?
    // futures::future::join_all(_x).await;

    receiver.iter().for_each(|(index, res)| {
        vec[index] = res;
    });

    println!("{:?}", vec);
}

async fn sleep_for(seconds: u64) {
    tokio::time::sleep(std::time::Duration::from_secs(seconds)).await;
}

async fn delay_thread_async_then_square(thread_no: usize, to_be_squared: usize) -> usize {
    let mut wait = 3;
    println!("Async thread {thread_no} sleeping for {wait} seconds");
    sleep_for(wait).await;
    wait = 7;
    println!("Async thread {thread_no} sleeping for {wait} seconds");
    sleep_for(wait).await;
    wait = 1;
    println!("Async thread {thread_no} sleeping for {wait} seconds");
    sleep_for(wait).await;
    wait = 4;
    println!("Async thread {thread_no} sleeping for {wait} seconds");
    sleep_for(wait).await;
    let res = to_be_squared.pow(2);
    println!("Async thread {thread_no} done. Result is {res}");
    res
}

已知修复方案

代码里有两处标记为FIX的修复方案:

  • 在main函数中调用with_tokio_channels_fail_call_fixing替代with_tokio_channels_fail_call_breaking
  • 在tokio_channels_fail函数中取消注释futures::future::join_all(_x).await

问题原因解答

核心错误点

你在with_tokio_channels_fail_call_breaking里犯了一个关键错误:在Tokio runtime之外的异步上下文里,用runtime的spawn提交任务后直接await这个handle,同时tokio_channels_fail内部又在执行阻塞式的receiver.iter().for_each调用,这直接导致了Tokio runtime的线程被耗尽,无法调度剩余任务。

1. 为什么部分任务能运行?

当你调用rt.spawn(tokio_channels_fail(...))时,Tokio runtime会把这个任务加入调度队列并开始执行。tokio_channels_fail内部会立刻创建并spawn出N个子任务,这些子任务一开始会被正常调度,直到它们遇到第一个tokio::time::sleep的await点——此时子任务会让出线程,回到runtime队列等待唤醒。

但在子任务全部进入await之前,tokio_channels_fail的代码已经走到了receiver.iter().for_each,这是一个完全阻塞的同步调用,会占用当前runtime线程的全部执行时间,不会给Tokio调度器让出线程资源。不过在阻塞发生前,已经有部分子任务完成了sleep前的代码(比如打印日志),甚至有些子任务可能在阻塞前就被唤醒完成了后续步骤,所以你能看到部分任务正常运行的迹象。

2. 为什么总是少执行一个任务?

这和std::sync::mpsc::channel的特性直接相关:只有当所有Sender实例都被drop后,receiver.iter()才会终止迭代。

你在代码中drop了主Sender,但每个子任务还持有一个克隆的Sender。当N-1个子任务完成后,它们持有的Sender会被自动drop,此时还剩最后一个Sender被未完成的子任务持有。但问题在于,这个最后一个子任务因为runtime线程被receiver.iter()阻塞,根本得不到调度机会,永远无法完成,它持有的Sender也不会被drop,导致receiver.iter()永远等待,程序无法退出——看起来就是“少执行了一个任务”。


修复方案生效原因

方案1:使用rt.block_on替代rt.spawn().await

rt.block_on是Tokio runtime提供的同步阻塞API,它会在runtime的上下文里执行传入的future。当future内部出现阻塞式调用时,Tokio的多线程runtime会自动把其他待调度任务转移到空闲线程上,不会让整个runtime陷入停滞。这样所有子任务都能获得调度机会,最终全部完成并drop各自的Sender,receiver.iter()就能正常终止。

方案2:异步等待所有子任务完成

取消注释futures::future::join_all(_x).await后,代码会异步等待所有子任务完成——这是一个非阻塞操作:当子任务进入await时,当前任务会让出线程,Tokio调度器可以继续调度其他任务。等所有子任务都完成后,它们的Sender会全部被drop,此时再执行receiver.iter()就能一次性拿到所有结果并正常终止,不会出现阻塞导致的调度问题。

内容的提问来源于stack exchange,提问作者NoBullsh1t

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.10 19:40:59