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

如何提升Rust Warp框架的API高并发吞吐量?

Warp API吞吐量远低于预期的排查与优化需求

我正在用Rust的Warp crate搭建首个API,但性能远未达到预期。为排除业务逻辑影响,我编写了一个极简的ping路由:

async fn ping() -> Result<impl Reply, Rejection> {
    Ok(html("pong"))
}

pub fn routes(app: Arc<App>) -> impl Filter<Extract = impl Reply, Error = Rejection> + Clone {
    path!("v1" / "pub" / "ping")
        .and(get())
        .and_then(ping)
}

用curl调用该端点结果正常,但用JMeter做压力测试时,最大吞吐量仅约2000请求/秒,且请求量超过数百后错误率高达60%以上,错误均为Connection refused。测试期间CPU使用率不足5%,表现远逊于TechEmpower基准测试宣称的50万请求/秒。

为排除JMeter的问题,我编写了一个Rust测试客户端,吞吐量提升至7000请求/秒,但仍不理想:

use core::fmt;
use std::error::Error;
use std::sync::Arc;
use std::time::Instant;

use hyper::{Request, StatusCode};
use hyper_util::rt::TokioIo;
use bytes::Bytes;
use http_body_util::Empty;
use tokio::net::TcpStream;
use tokio::task;

type Result<T> = std::result::Result<T, Box<dyn Error + Send + Sync>>;

#[tokio::main]
async fn main() -> Result<()> {
    let url = "http://localhost:8000/v1/pub/ping".parse::<hyper::Uri>()?;
    let concurrency = 20;
    let repeat_count = 500;

    let authority = url.authority().unwrap();

    let addr = Arc::new(String::from(authority.as_str()));
    let req = Arc::new(Request::builder()
        .method("GET")
        .uri(url.path())
        .header(hyper::header::HOST, authority.as_str())
        .header(hyper::header::CONNECTION, "keep-alive")
        .body(Empty::<Bytes>::new())?);

    // Warm up the connection
    fetch(addr.clone(), req.clone()).await?;

    let start = Instant::now();

    let mut tasks = Vec::new();

    for _ in 0..concurrency {
        tasks.push(task::spawn(fetch(addr.clone(), req.clone())));
    }

    for _ in 1..repeat_count {
        for _ in 0..concurrency {
            let task = tasks.remove(0);
            let _ = task.await?;
            tasks.push(task::spawn(fetch(addr.clone(), req.clone())));
        }
    }

    for task in tasks {
        let _ = task.await?;
    }

    let elapsed = start.elapsed();
    println!("Elapsed: {:.2?} GET {} {} times with {} concurrency", elapsed, url, concurrency * repeat_count, concurrency);
    println!("Average {} req/sec", 1000 * concurrency * repeat_count / elapsed.as_millis());

    Ok(())
}

async fn fetch(addr: Arc<String>, req: Arc<Request<Empty<Bytes>>>) -> Result<()> {
    let stream = TcpStream::connect(&*addr).await?;
    let io = TokioIo::new(stream);
    
    let (mut sender, conn) = hyper::client::conn::http1::handshake(io).await?;
    task::spawn(async move {
        if let Err(err) = conn.await {
            println!("Connection failed: {:?}", err);
        }
    });

    let res = sender.send_request((*req).clone()).await?;

    match res.status() {
        StatusCode::OK => Ok(()),
        _ => Err(Box::new(HttpError::new(res.status())))
    }
}

#[derive(Debug)]
struct HttpError {
    pub status: StatusCode,
}

impl HttpError {
    pub fn new(status: StatusCode) -> Self {
        Self { status }
    }
}

impl fmt::Display for HttpError {
    fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
        write!(f, "Http status code {}", self.status )
    }
}

impl Error for HttpError {
    fn source(&self) -> Option<&(dyn Error + 'static)> { None }
    fn description(&self) -> &str { "" }
    fn cause(&self) -> Option<&dyn Error> { self.source() }
}

随后我基于标准库编写了一个极简Web服务器,吞吐量约7500请求/秒,略高于Warp的表现:

use crate::App;
use std::{
    io::{prelude::*, BufReader},
    net::{Ipv4Addr, TcpListener, TcpStream},
    sync::{atomic::Ordering, mpsc, Arc},
    thread::{self, available_parallelism, Thread},
};

pub fn run(app: Arc<App>, ipv4: Ipv4Addr, port: u16) {
    let authority = format!("{}:{}", ipv4, port);
    let listener = TcpListener::bind(authority).unwrap();

    let concurrency = available_parallelism().unwrap().get();
    let mut threads: Vec<thread::JoinHandle<()>> = Vec::with_capacity(concurrency);
    let mut senders: Vec<mpsc::Sender<TcpStream>> = Vec::with_capacity(concurrency);

    for _ in 0..concurrency {
        let (sender, receiver) = mpsc::channel();
        senders.push(sender);

        let app = app.clone();
        threads.push(thread::spawn(move || process_connections(app, receiver)))
    }

    let mut thread_index = 0;
    for stream in listener.incoming() {
        senders[thread_index].send(stream.unwrap()).unwrap();
        thread_index = (thread_index + 1) % concurrency;
    }
}

fn process_connections(app: Arc<App>, receiver: mpsc::Receiver<TcpStream>) {
    loop {
        let stream = receiver.recv().unwrap();
        handle_connection(app.clone(), stream);
    }
}

fn handle_connection(app: Arc<App>, mut stream: TcpStream) {
    let buf_reader = BufReader::new(&stream);
    let _http_request: Vec<_> = buf_reader
        .lines()
        .map(|result| result.unwrap())
        .take_while(|line| !line.is_empty())
        .collect();

    let response = "HTTP/1.1 200 OK

";
    stream.write_all(response.as_bytes()).unwrap();

    app.clone()
        .request_count
        .clone()
        .fetch_add(1, Ordering::Relaxed);
}

我怀疑问题出在客户端实现,或是MacBook Pro的系统固有限制,打算尝试单连接多请求的流式传输方案,现寻求能有效提升Warp吞吐量的具体方案。

内容的提问来源于stack exchange,提问作者bikeman868

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 21:44:56