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

Rust线程池异常:同一Worker处理多请求问题排查求助

多线程Web服务器线程池调度异常问题

我正在把一个带错误处理、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),
        }
    }
}

问题根源分析

  1. Handler锁持有时间过长:任务执行时,handler.lock().unwrap()会获取Mutex锁,且在整个handle_request/handle_bad_request执行期间持续持有。当Worker处理/sleep请求时,锁会被持有5秒,其他Worker即使拿到任务也无法获取锁执行逻辑,只能等待,导致看起来像同一个Worker处理两个请求。

  2. 线程池任务调度锁竞争: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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 16:24:54