如何提升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
相关产品推荐
相关产品推荐

