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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 18:06:01