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

基于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(信号量)**,它可以直接控制同时执行的任务数量,无需手动维护任务池和通道。

实现思路

  1. 初始化一个信号量,许可数量等于你要限制的并发数(比如20)。
  2. 每次处理请求时,先从信号量获取一个许可(无许可则等待)。
  3. 在异步任务中执行外部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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.26 01:18:27