为何以异常方式调用Tokio时仅运行n-1个任务?
我最近在学习Rust并发/并行编程,编写不同并发实现的代码片段时碰到了Tokio的异常问题。我确认代码是误用了Tokio,程序没崩溃但运行不正常——总是少执行一个任务。我已经找到规避方法,但想搞懂内部原因,有两个疑问:
- 既然代码确实写错了,为什么还有部分任务能正常运行?
- 为什么偏偏总是少执行一个任务?
复现代码
#![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

