如何在Rust中用队列处理异步任务?解决API并发问题
解决异步任务串行队列处理问题
你当前的代码每次请求都会通过tokio::spawn启动新任务,这些任务会在Tokio的线程池中并行执行,导致资源消耗超出预期。要实现任务按队列顺序逐个处理,可以用Tokio的mpsc通道构建一个串行任务队列,确保任务依次执行。
修改后的完整代码
use axum::{extract::Path, routing::get, Router}; use tokio::{sync::mpsc, time::Duration}; use std::sync::Arc; extern crate diesel; extern crate tracing; #[tokio::main] async fn main() { tracing_subscriber::fmt::init(); // 创建mpsc通道,缓冲区大小可根据业务需求调整 let (sender, mut receiver) = mpsc::channel(100); // 用Arc包裹Sender,实现多请求间的线程安全共享 let shared_sender = Arc::new(sender); // 启动后台串行处理任务:单个消费者循环接收并执行任务 tokio::spawn(async move { while let Some(timer) = receiver.recv().await { start_timer_send_json(timer).await; } }); let app = Router::new() .route("/sleep/:id", get(sleep_and_print)) // 将共享的Sender注入路由状态,供请求处理函数使用 .with_state(shared_sender); let addr = std::net::SocketAddr::from(([0, 0, 0, 0], 3000)); tracing::info!("Listening on {}", addr); axum::Server::bind(&addr) .serve(app.into_make_service()) .await .unwrap(); } // 从路由状态中获取Sender,将任务发送到串行队列 async fn sleep_and_print(Path(timer): Path<i32>, sender: Arc<mpsc::Sender<i32>>) -> String { // 处理任务发送失败的情况(比如队列缓冲区已满) if sender.send(timer).await.is_err() { return "{\"error\": \"Failed to queue task\"}".to_string(); } format!("{{\"timer\": {}}}", timer) } async fn start_timer_send_json(timer: i32) { println!("Start timer {}.", timer); tokio::time::sleep(Duration::from_secs(300)).await; println!("Timer {} done.", timer); }
核心逻辑说明
- mpsc通道:
mpsc::channel创建多生产者单消费者通道,所有API请求作为生产者将任务发送到队列,后台单个消费者任务依次接收并执行,保证任务串行处理。 - Arc共享Sender:用
Arc包裹Sender,实现多个请求线程对发送器的安全共享访问。 - 后台串行任务:主线程启动的单独spawn任务会持续监听队列,一旦有新任务就执行,确保所有任务按提交顺序逐个运行,避免资源并行消耗。
内容的提问来源于stack exchange,提问作者0xRoronoa Zoro
相关产品推荐
相关产品推荐

