Rust线程池异常:同一Worker处理多请求问题排查求助
我正在把一个带错误处理、Trait扩展的自定义单线程Web服务器改成多线程版本。已经完成了手册里的线程池示例,自己实现的线程池能正常编译,但测试时发现,当同时请求/sleep路径(会阻塞5秒)和普通路径时,同一个Worker会处理这两个请求,而不是在第一个请求阻塞期间让其他Worker处理第二个请求,求解决思路。
服务器run方法实现
pub fn run(self, handler: impl Handler) { println!("Listening on {}", self.addr); let listener = TcpListener::bind(&self.addr).unwrap(); let pool = ThreadPool::new(4); let handler = Arc::new(Mutex::new(handler)); for stream in listener.incoming() { let mut stream = stream.unwrap(); let mut buffer = [0; 1024]; let handler = handler.clone(); pool.execute(move || match stream.read(&mut buffer) { Ok(_) => { println!("{}", String::from_utf8_lossy(&buffer)); let response = match Request::try_from(&buffer[..]) { Ok(request) => handler.lock().unwrap().handle_request(&request), Err(e) => handler.lock().unwrap().handle_bad_request(&e), }; if let Err(e) = response.send(&mut stream) { println!("Failed to send response: {}", e); } } Err(e) => println!("Failed to read from connection: {}", e), }); } }
ThreadPool实现
use std::{ sync::{mpsc, Arc, Mutex}, thread, }; pub struct ThreadPool { workers: Vec<Worker>, sender: mpsc::Sender<Job>, } impl ThreadPool { pub fn new(size: usize) -> ThreadPool { assert!(size > 0); let (sender, receiver) = mpsc::channel(); let receiver = Arc::new(Mutex::new(receiver)); let mut workers = Vec::with_capacity(size); for id in 0..size { workers.push(Worker::new(id, Arc::clone(&receiver))); } ThreadPool { workers, sender } } pub fn execute<F>(&self, f: F) where F: FnOnce() + Send + 'static, { let job = Box::new(f); self.sender.send(job).unwrap(); } } type Job = Box<dyn FnOnce() + Send + 'static>; struct Worker { id: usize, thread: thread::JoinHandle<()>, } impl Worker { fn new(id: usize, receiver: Arc<Mutex<mpsc::Receiver<Job>>>) -> Worker { let thread = thread::spawn(move || loop { let job = receiver.lock().unwrap().recv().unwrap(); println!("Worker {} got a job: executing.", id); job(); }); Worker { id, thread } } }
Handler Trait定义
pub trait Handler: Send + 'static { fn handle_request(&mut self, request: &Request) -> Response; fn handle_bad_request(&mut self, e: &ParseError) -> Response { println!("Failed to parse request: {}", e); Response::new(StatusCode::BadRequest, None) } }
WebsiteHandler实现Handler
impl Handler for WebsiteHandler { fn handle_request(&mut self, request: &Request) -> Response { match request.method() { Method::GET => match request.path() { "/" => Response::new(StatusCode::Ok, self.read_file("index.html")), "/hello" => Response::new(StatusCode::Ok, self.read_file("hello.html")), "/sleep" => { thread::sleep(Duration::from_secs(5)); Response::new(StatusCode::Ok, self.read_file("hello.html")) } path => match self.read_file(path) { Some(contents) => Response::new(StatusCode::Ok, Some(contents)), None => Response::new(StatusCode::NotFound, None), }, }, _ => Response::new(StatusCode::NotFound, None), } } }
问题根源分析
Handler锁持有时间过长:任务执行时,
handler.lock().unwrap()会获取Mutex锁,且在整个handle_request/handle_bad_request执行期间持续持有。当Worker处理/sleep请求时,锁会被持有5秒,其他Worker即使拿到任务也无法获取锁执行逻辑,只能等待,导致看起来像同一个Worker处理两个请求。线程池任务调度锁竞争:Worker循环中,
receiver.lock().unwrap().recv()会先获取接收端Mutex锁再调用recv,而recv是阻塞的——这意味着Worker在等待任务时会一直持有接收端锁,其他Worker无法抢锁接收新任务,只能等当前Worker释放锁后才能分配任务,导致任务串行处理。
解决方案
1. 缩短Handler锁的持有时间(最优解)
如果WebsiteHandler的read_file不需要修改内部状态,直接将Handler的方法签名改为&self,用Arc代替Arc<Mutex>,彻底消除锁竞争:
- 修改Handler Trait:
pub trait Handler: Send + Sync + 'static { fn handle_request(&self, request: &Request) -> Response; fn handle_bad_request(&self, e: &ParseError) -> Response { println!("Failed to parse request: {}", e); Response::new(StatusCode::BadRequest, None) } }
- 修改WebsiteHandler实现:
impl Handler for WebsiteHandler { fn handle_request(&self, request: &Request) -> Response { // 原有逻辑不变,仅将&mut self改为&self match request.method() { Method::GET => match request.path() { "/" => Response::new(StatusCode::Ok, self.read_file("index.html")), "/hello" => Response::new(StatusCode::Ok, self.read_file("hello.html")), "/sleep" => { thread::sleep(Duration::from_secs(5)); Response::new(StatusCode::Ok, self.read_file("hello.html")) } path => match self.read_file(path) { Some(contents) => Response::new(StatusCode::Ok, Some(contents)), None => Response::new(StatusCode::NotFound, None), }, }, _ => Response::new(StatusCode::NotFound, None), } } }
- 修改run方法中Handler的包装:
let handler = Arc::new(handler); // 移除Mutex包装
2. 修复线程池任务调度逻辑
缩小接收端锁的持有范围,仅在获取任务瞬间持有锁:
修改Worker的new方法:
impl Worker { fn new(id: usize, receiver: Arc<Mutex<mpsc::Receiver<Job>>>) -> Worker { let thread = thread::spawn(move || loop { // 通过作用域限制锁的持有时间,拿到任务后立即释放 let job = { let mut receiver = receiver.lock().unwrap(); receiver.recv().unwrap() }; println!("Worker {} got a job: executing.", id); job(); }); Worker { id, thread } } }
3. 验证效果
修改后,处理/sleep请求的Worker会阻塞5秒,但其他Worker可以同时接收并处理新请求,真正实现多线程并行处理。
内容的提问来源于stack exchange,提问作者Win Barua

