基于crossbeam channel的异步请求池阻塞问题排查及优化咨询
问题描述
我计划开发一个API服务器,需从外部API获取部分信息,但该外部API稳定性较差,因此希望将并发请求数限制在10或20以内。
我尝试实现一个HttpPool,通过crossbeam有界通道接收任务并分配给tokio任务,以此控制并发请求量。但发现当任务数超过8个时,程序在获取首个任务后就会阻塞。
以下是我的Rust代码:
use std::{error::Error, result::Result}; use tokio::sync::oneshot::Sender; use tokio::time::timeout; use tokio::time::{sleep, Duration}; use crossbeam_channel; #[derive(Debug)] struct HttpTaskRequest { url: String, result: Sender<String>, } type PoolSender = crossbeam_channel::Sender<HttpTaskRequest>; type PoolReceiver = crossbeam_channel::Receiver<HttpTaskRequest>; #[derive(Debug)] struct HttpPool { size: i32, sender: PoolSender, receiver: PoolReceiver, } impl HttpPool { fn new(capacity: i32) -> Self { let (tx, rx) = crossbeam_channel::bounded::<HttpTaskRequest>(capacity as usize); HttpPool { size: capacity, sender: tx, receiver: rx, } } async fn start(self) -> Result<HttpPool, Box<dyn Error>> { for i in 0..self.size { let task_receiver = self.receiver.clone(); tokio::spawn(async move { loop { match task_receiver.recv() { Ok(request) => { if request.result.is_closed() { println!("Task[{i}] received url {} already closed by receiver, seems to reach timeout already", request.url); } else { println!("Task[{i}] started to work {:?}", request.url); let resp = reqwest::get("https://httpbin.org/ip").await; println!("Resp: {:?}", resp); println!("Done Send request for url {}", request.url); request.result.send("Result".to_owned()).expect("Failed to send result"); } } Err(err) => println!("Error: {err}"), } } }); } Ok(self) } pub async fn request(&self, url: String) -> Result<(), Box<dyn Error>> { let (os_sender, os_receiver) = tokio::sync::oneshot::channel::<String>(); let request = HttpTaskRequest { result: os_sender, url: url.clone(), }; self.sender.send(request).expect("Failed to publish message to task group"); match timeout(Duration::from_millis(100), os_receiver).await { Ok(res) => { println!("Request finished without reaching the timeout {}",res.unwrap()); } Err(_) => {println!("Request {url} run into timeout");} } Ok(()) } } #[tokio::main] async fn main() { let http_pool = HttpPool::new(20).start().await.expect("Failed to start http pool"); for i in 0..10 { let url = format!("T{}", i.to_string()); http_pool.request(url).await.expect("Failed to request message"); } loop {} }
请问该代码阻塞的原因是什么?是否与tokio::spawn有关?同时希望了解实现动态并发请求限制的正确方案。我曾参考过相关问题,但该方案适用于已知任务量的场景,无法满足我这里任务动态到来的需求。
阻塞原因分析
代码阻塞的核心问题不是tokio::spawn,而是错误地在异步任务中使用了同步阻塞的通道接收操作:
crossbeam_channel::recv()是同步阻塞方法,调用后会直接占用当前tokio工作线程,直到有消息到达才释放。- tokio默认工作线程数等于CPU核心数(比如8核机器就是8个线程),当你启动超过8个
tokio::spawn任务后,每个任务都在loop里调用recv()阻塞线程,会把所有工作线程占满。 - 后续的异步任务(包括请求响应逻辑)无法获得工作线程调度,导致整个程序卡住。
动态并发限制的正确实现方案
对于动态到来的任务,最简洁且符合tokio异步模型的方案是使用**tokio::sync::Semaphore(信号量)**,它可以直接控制同时执行的任务数量,无需手动维护任务池和通道。
实现思路
- 初始化一个信号量,许可数量等于你要限制的并发数(比如20)。
- 每次处理请求时,先从信号量获取一个许可(无许可则等待)。
- 在异步任务中执行外部API请求,完成后释放许可,让后续任务继续执行。
示例代码
use std::{error::Error, sync::Arc}; use tokio::sync::Semaphore; use tokio::time::{timeout, Duration}; struct HttpPool { semaphore: Arc<Semaphore>, } impl HttpPool { fn new(max_concurrent: usize) -> Self { HttpPool { semaphore: Arc::new(Semaphore::new(max_concurrent)), } } pub async fn request(&self, url: String) -> Result<String, Box<dyn Error>> { // 获取许可,没有则等待 let permit = self.semaphore.acquire().await?; // 执行外部API请求,添加超时控制 let result = timeout(Duration::from_millis(100), async { let resp = reqwest::get("https://httpbin.org/ip").await?; let text = resp.text().await?; Ok(text) }).await; // 释放许可(permit销毁时会自动释放,这里手动drop更直观) drop(permit); match result { Ok(Ok(text)) => { println!("Request {} finished: {}", url, text); Ok(text) } Ok(Err(e)) => Err(e.into()), Err(_) => { println!("Request {} timed out", url); Err("Request timed out".into()) } } } } #[tokio::main] async fn main() -> Result<(), Box<dyn Error>> { let http_pool = HttpPool::new(20); // 模拟动态到来的10个请求 for i in 0..10 { let url = format!("T{}", i); tokio::spawn({ let pool = http_pool.clone(); async move { let _ = pool.request(url).await; } }); } // 监听中断信号,保持程序运行 tokio::signal::ctrl_c().await?; Ok(()) }
方案优势
- 完全基于tokio异步模型,不会出现同步阻塞占用线程的问题。
- 自动处理并发限制,无需手动维护任务队列和工作线程。
- 支持动态任务:无论任务是批量还是逐个到来,信号量都会自动控制并发数,无需提前知道任务总量。
内容的提问来源于stack exchange,提问作者whati001
相关产品推荐
相关产品推荐

