使用tokio::time::timeout判断超时后程序挂起问题咨询
问题描述
如下代码的输入是简单csv文件,格式为每行两个字段,分别对应问题和答案:
#[derive(Debug)] struct Problem { q: String, a: String, } impl Problem { pub fn new(q: &str, a: &str) -> Problem { Self { q: q.to_owned(), a: a.to_owned(), } } } const FILE_NAME: &str = "../../input/problems.csv"; #[tokio::main] async fn main() { let probs = parse_lines().expect("Could not parse CSV file"); println!("Banana quiz about to start. Press enter when ready."); let mut buf = String::new(); match std::io::stdin().read_line(&mut buf).ok() { None => { println!("Error reading user input"); std::process::exit(1) } _ => {} } let mut correct_ans = 0; for p in &probs { println!("What banana? {}", p.q); let (banana_s, mut banana_r) = mpsc::channel(512); tokio::spawn(async move { let mut banana = String::new(); std::io::stdin().read_line(&mut banana); banana.pop(); banana_s.send(banana).await.unwrap(); }); let mut banana = String::new(); match tokio::time::timeout(Duration::from_secs(5), banana_r.recv()).await { Ok(opt) => { match opt { Some(b) => banana.push_str(&b), None => {}, } }, Err (_) => { println!("Only have 5 seconds to input the answer!"); return; } }; println!("Your Banana: {}, Correct Banana: {}\n", banana, p.a); if banana != p.a { println!("BAD BANANA!"); break; } println!("good banana!"); println!("------------"); correct_ans += 1; } println!("Correct answers: {}/{}\n", correct_ans, probs.len()) } fn parse_lines() -> Result<Vec<Problem>, csv::Error> { let mut builder = ReaderBuilder::new(); builder.has_headers(false); let mut reader = builder.from_path(FILE_NAME)?; let mut probs = Vec::new(); for r in reader.records() { let rec = r?; probs.push(Problem::new(&rec[0], &rec[1])); } Ok(probs) }
除用户输入答案超过5秒的场景外,其余逻辑运行正常。打印Only have 5 seconds to input the answer提示后,程序不会直接退出,必须按下回车才会继续执行,且回车后会触发如下panic:
thread 'tokio-runtime-worker' panicked at 'called `Result::unwrap()` on an `Err` value: SendError("")', src/main.rs:46:41
报错指向通道发送数据的代码行:
banana_s.send(banana).await.unwrap();
推测该问题是tokio spawn的任务被read_line()函数阻塞导致,而使用std::thread搭配std::sync::mpsc::channel时不存在该问题,程序可正常退出。问题核心为两点:该问题的根本原因是什么?如何实现超时触发后程序直接正常退出?
解答
根本原因
- 你使用了标准库的阻塞式stdin而非Tokio提供的异步stdin,
std::io::stdin().read_line()是阻塞系统调用,会直接卡住当前执行它的Tokio工作线程,直到用户输入回车才会返回。超时触发时该任务还卡在read_line调用上,根本没有执行到后续send逻辑,所以程序不会立刻退出。 - 超时触发后,main函数中的接收端
banana_r会随着函数return被销毁,Tokio MPSC通道只要接收端销毁,发送端调用send就会返回SendError。等用户按下回车,被卡住的任务终于恢复执行,调用send时就会因为接收端已不存在触发unwrap panic。 - 用
std::thread没有该问题的原因是:std::thread是独立的系统线程,阻塞后不会影响主线程执行,超时后主线程直接退出整个进程,所有线程都会被系统强制回收,自然不会出现后续send触发panic的问题。而当前代码超时后仅从main的async块返回,Tokio runtime默认会等待所有spawn的任务执行完成才会完全退出,所以会卡住等待阻塞的read_line返回。
修复方案
方案1:超时后直接终止进程
最简单的修改方式,把超时分支里的return替换为std::process::exit(1),直接终止整个进程,不管还有没有未执行完的任务,既不会等待用户输入,也不会出现后续panic:
Err (_) => { println!("Only have 5 seconds to input the answer!"); std::process::exit(1); // 替换原有return逻辑 }
方案2:使用Tokio异步stdin替换阻塞式stdin
把标准库stdin换成Tokio提供的异步实现,不会阻塞工作线程,超时后任务可以被正常取消,同时去掉send的unwrap避免接收端销毁时panic:
// 首先导入异步读扩展:use tokio::io::AsyncBufReadExt; tokio::spawn(async move { let stdin = tokio::io::stdin(); let mut reader = tokio::io::BufReader::new(stdin); let mut banana = String::new(); let _ = reader.read_line(&mut banana).await; banana.pop(); // 忽略send错误,接收端销毁时直接丢弃数据即可 let _ = banana_s.send(banana).await; });
方案3:阻塞读放到spawn_blocking中执行
如果必须使用标准库的阻塞stdin,就放到Tokio专门为阻塞任务设计的spawn_blocking接口中执行,不会占用异步工作线程:
tokio::task::spawn_blocking(move || { let mut banana = String::new(); let _ = std::io::stdin().read_line(&mut banana); banana.pop(); // 阻塞任务中执行异步send tokio::runtime::Handle::current().block_on(async { let _ = banana_s.send(banana).await; }); });
内容的提问来源于stack exchange,提问作者antogilbert
相关产品推荐
相关产品推荐

