基于Tokio的Rust程序输出结果后挂起未正常退出问题排查
问题描述
用Rust编写了一个基础HackerNews客户端,通过其API获取前500条热门故事的标题。程序能快速输出所有标题,但之后会挂起10-15秒才退出。怀疑是Tokio运行时的任务清理导致此问题,希望了解原因及可行的优化方案,不愿使用std::process::exit(0)这类强制退出的方式。
最小复现代码
src/main.rs:
use futures::future::join_all; const BASE_URL: &str = "https://hacker-news.firebaseio.com/v0"; #[tokio::main(flavor = "multi_thread", worker_threads = 10)] async fn main() { let client = reqwest::Client::new(); // 获取所有热门故事的ID let top_stories_ids: Vec<i32> = client .get(format!("{}/topstories.json", BASE_URL)) .send() .await .unwrap() .json::<Vec<i32>>() .await .unwrap(); // 生成获取每个故事详情的Future列表 let top_stories_items_futures: Vec<_> = top_stories_ids .into_iter() .map(|id| client.get(format!("{}/item/{}.json", BASE_URL, id)).send()) .collect(); // 批量发起网络请求,生成解析JSON的Future列表 let top_stories_json_futures: Vec<_> = join_all(top_stories_items_futures) .await .into_iter() .map(|resp| resp.unwrap().json::<serde_json::Value>()) .collect(); // 批量解析JSON并收集结果,打印所有标题 let top_stories_items: Vec<_> = join_all(top_stories_json_futures) .await .into_iter() .filter_map(|json| json.ok()) .collect(); for story in &top_stories_items { println!("{:?}", story["title"]); } }
Cargo依赖
futures = "0.3.28" reqwest = { version = "0.11.18", features = ["json", "rustls-tls"] } serde = { version = "1.0.166", features = ["derive"] } serde_json = "1.0.100" tokio = { version = "1.29.1", features = ["macros", "rt-multi-thread"] }
原因分析
问题核心在于reqwest连接池的默认行为与Tokio运行时的清理逻辑:
- reqwest的
Client默认会创建持久化连接池,空闲连接会保持90秒的存活时间以便复用。 - 一次性发起500个请求会生成大量连接,程序完成任务后,Tokio运行时需要等待这些连接池内的后台任务(如连接超时清理)自然结束才会完全关闭,导致程序挂起。
优化方案
方案1:显式销毁reqwest Client
在程序结束前主动丢弃Client,触发连接池的销毁逻辑,让所有空闲连接立即断开:
// 在打印标题后添加 drop(client); // 可选:等待100毫秒确保连接完全关闭(多数场景下无需此步骤) tokio::time::sleep(tokio::time::Duration::from_millis(100)).await;
方案2:缩短连接池空闲超时时间
创建Client时通过ClientBuilder设置更短的空闲连接超时,让连接更快被回收:
let client = reqwest::ClientBuilder::new() .pool_idle_timeout(std::time::Duration::from_secs(1)) // 设置1秒空闲超时 .build() .unwrap();
方案3:限制并发请求数量
一次性发起500个请求会导致连接池压力过大,同时可能触发API限流。用FuturesUnordered控制并发数,减少同时存在的连接:
use futures::stream::{FuturesUnordered, StreamExt}; // 替换原有join_all的逻辑 let mut stream = FuturesUnordered::new(); const MAX_CONCURRENT: usize = 20; // 限制并发数为20 let mut in_flight = 0; for id in top_stories_ids { if in_flight >= MAX_CONCURRENT { // 等待一个请求完成后再添加新请求 if let Some(_) = stream.next().await { in_flight -= 1; } } stream.push(client.get(format!("{}/item/{}.json", BASE_URL, id)).send()); in_flight += 1; } // 处理剩余未完成的请求 while let Some(resp) = stream.next().await { match resp { Ok(r) => { if let Ok(json) = r.json::<serde_json::Value>().await { if let Some(title) = json.get("title") { println!("{:?}", title); } } } Err(e) => eprintln!("请求失败: {}", e), } }
内容的提问来源于stack exchange,提问作者AUD_FOR_IUV
相关产品推荐
相关产品推荐

