Keyboard shortcuts

Press or to navigate between chapters

Press S or / to search in the book

Press ? to show this help

Press Esc to hide this help

Корректное завершение и очистка

Код в листинге 21-20 отвечает на запросы асинхронно с помощью пула потоков, как мы и планировали. Мы получаем несколько предупреждений о полях workers, id и thread, которые не используем напрямую; они напоминают нам, что мы ничего не очищаем. Когда мы используем менее изящный метод ctrl-C, чтобы остановить основной поток, все остальные потоки тоже останавливаются немедленно, даже если они находятся в середине обслуживания запроса.

Далее мы реализуем трейт Drop, чтобы вызывать join для каждого потока в пуле, позволяя им завершить запросы, над которыми они работают, прежде чем закрыться. Затем мы реализуем способ сообщать потокам, что им нужно перестать принимать новые запросы и завершить работу. Чтобы увидеть этот код в действии, мы изменим сервер так, чтобы он принимал только два запроса перед корректным завершением своего пула потоков.

Стоит заметить по ходу дела: ничто из этого не влияет на части кода, которые отвечают за выполнение замыканий, поэтому все здесь было бы таким же, если бы мы использовали пул потоков для async runtime.

Реализация трейта Drop для ThreadPool

Начнем с реализации Drop для нашего пула потоков. Когда пул сбрасывается, все наши потоки должны присоединиться, чтобы гарантировать завершение своей работы. Листинг 21-22 показывает первую попытку реализации Drop; этот код пока не совсем будет работать.

Имя файла: src/lib.rs
use std::{
    sync::{Arc, Mutex, mpsc},
    thread,
};

pub struct ThreadPool {
    workers: Vec<Worker>,
    sender: mpsc::Sender<Job>,
}

type Job = Box<dyn FnOnce() + Send + 'static>;

impl ThreadPool {
    /// Create a new ThreadPool.
    ///
    /// The size is the number of threads in the pool.
    ///
    /// # Panics
    ///
    /// The `new` function will panic if the size is zero.
    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();
    }
}

impl Drop for ThreadPool {
    fn drop(&mut self) {
        for worker in &mut self.workers {
            println!("Shutting down worker {}", worker.id);

            worker.thread.join().unwrap();
        }
    }
}

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 {id} got a job; executing.");

                job();
            }
        });

        Worker { id, thread }
    }
}
Listing 21-22: Присоединение каждого потока при выходе пула из области видимости

Сначала мы проходим циклом по всем workers пула потоков. Мы используем &mut, потому что self является изменяемой ссылкой, и нам также нужно иметь возможность изменять worker. Для каждого worker мы печатаем сообщение о том, что конкретный экземпляр Worker завершает работу, а затем вызываем join у потока этого экземпляра Worker. Если вызов join завершится неудачно, мы используем unwrap, чтобы Rust запаниковал и перешел к некорректному завершению.

Вот ошибка, которую мы получаем при компиляции этого кода:

$ cargo check
    Checking hello v0.1.0 (file:///projects/hello)
error[E0507]: cannot move out of `worker.thread` which is behind a mutable reference
  --> src/lib.rs:52:13
   |
52 |             worker.thread.join().unwrap();
   |             ^^^^^^^^^^^^^ ------ `worker.thread` moved due to this method call
   |             |
   |             move occurs because `worker.thread` has type `JoinHandle<()>`, which does not implement the `Copy` trait
   |
note: `JoinHandle::<T>::join` takes ownership of the receiver `self`, which moves `worker.thread`
  --> /rustc/1159e78c4747b02ef996e55082b704c09b970588/library/std/src/thread/mod.rs:1921:17

For more information about this error, try `rustc --explain E0507`.
error: could not compile `hello` (lib) due to 1 previous error

Ошибка говорит нам, что мы не можем вызвать join, потому что у нас есть только изменяемое заимствование каждого worker, а join забирает владение своим аргументом. Чтобы решить эту проблему, нужно переместить поток из экземпляра Worker, владеющего thread, чтобы join мог потребить поток. Один способ сделать это — использовать тот же подход, что и в листинге 18-15. Если бы Worker хранил Option<thread::JoinHandle<()>>, мы могли бы вызвать метод take у Option, чтобы переместить значение из варианта Some и оставить на его месте вариант None. Другими словами, у работающего Worker в thread был бы вариант Some, а когда мы захотели бы очистить Worker, мы заменили бы Some на None, чтобы у Worker больше не было потока для запуска.

Однако единственный момент, когда это понадобилось бы, — сброс Worker. В обмен на это нам пришлось бы работать с Option<thread::JoinHandle<()>> везде, где мы обращаемся к worker.thread. Идиоматичный Rust довольно часто использует Option, но когда вы ловите себя на том, что оборачиваете в Option что-то, о чем знаете, что оно всегда будет присутствовать, как обходной путь вроде этого, стоит поискать альтернативные подходы, чтобы сделать код чище и менее склонным к ошибкам.

В этом случае существует лучшая альтернатива: метод Vec::drain. Он принимает параметр диапазона, чтобы указать, какие элементы удалить из вектора, и возвращает итератор этих элементов. Передача синтаксиса диапазона .. удалит каждое значение из вектора.

Поэтому нам нужно обновить реализацию drop для ThreadPool вот так:

Имя файла: src/lib.rs
#![allow(unused)]
fn main() {
use std::{
    sync::{Arc, Mutex, mpsc},
    thread,
};

pub struct ThreadPool {
    workers: Vec<Worker>,
    sender: mpsc::Sender<Job>,
}

type Job = Box<dyn FnOnce() + Send + 'static>;

impl ThreadPool {
    /// Create a new ThreadPool.
    ///
    /// The size is the number of threads in the pool.
    ///
    /// # Panics
    ///
    /// The `new` function will panic if the size is zero.
    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();
    }
}

impl Drop for ThreadPool {
    fn drop(&mut self) {
        for worker in self.workers.drain(..) {
            println!("Shutting down worker {}", worker.id);

            worker.thread.join().unwrap();
        }
    }
}

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 {id} got a job; executing.");

                job();
            }
        });

        Worker { id, thread }
    }
}
}

Это устраняет ошибку компилятора и не требует никаких других изменений в нашем коде. Обратите внимание: поскольку drop может быть вызван во время паники, unwrap также может запаниковать и вызвать двойную панику, которая немедленно аварийно завершит программу и остановит любую выполняющуюся очистку. Для примерной программы это допустимо, но для производственного кода не рекомендуется.

Сигнал потокам прекратить ожидание заданий

Со всеми внесенными изменениями наш код компилируется без предупреждений. Однако плохая новость в том, что этот код пока работает не так, как мы хотим. Ключ находится в логике замыканий, выполняемых потоками экземпляров Worker: сейчас мы вызываем join, но это не завершит потоки, потому что они вечно выполняют loop в поисках заданий. Если мы попробуем сбросить наш ThreadPool с текущей реализацией drop, основной поток заблокируется навсегда, ожидая завершения первого потока.

Чтобы исправить эту проблему, нам понадобится изменение в реализации drop для ThreadPool, а затем изменение в цикле Worker.

Сначала мы изменим реализацию drop для ThreadPool, чтобы явно сбросить sender перед ожиданием завершения потоков. Листинг 21-23 показывает изменения в ThreadPool для явного сброса sender. В отличие от потока, здесь нам действительно нужно использовать Option, чтобы иметь возможность переместить sender из ThreadPool с помощью Option::take.

Имя файла: src/lib.rs
use std::{
    sync::{Arc, Mutex, mpsc},
    thread,
};

pub struct ThreadPool {
    workers: Vec<Worker>,
    sender: Option<mpsc::Sender<Job>>,
}
// --snip--

type Job = Box<dyn FnOnce() + Send + 'static>;

impl ThreadPool {
    /// Create a new ThreadPool.
    ///
    /// The size is the number of threads in the pool.
    ///
    /// # Panics
    ///
    /// The `new` function will panic if the size is zero.
    pub fn new(size: usize) -> ThreadPool {
        // --snip--

        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: Some(sender),
        }
    }

    pub fn execute<F>(&self, f: F)
    where
        F: FnOnce() + Send + 'static,
    {
        let job = Box::new(f);

        self.sender.as_ref().unwrap().send(job).unwrap();
    }
}

impl Drop for ThreadPool {
    fn drop(&mut self) {
        drop(self.sender.take());

        for worker in self.workers.drain(..) {
            println!("Shutting down worker {}", worker.id);

            worker.thread.join().unwrap();
        }
    }
}

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 {id} got a job; executing.");

                job();
            }
        });

        Worker { id, thread }
    }
}
Listing 21-23: Явный сброс sender перед присоединением потоков Worker

Сброс sender закрывает канал, что указывает: больше сообщений отправляться не будет. Когда это происходит, все вызовы recv, которые экземпляры Worker делают в бесконечном цикле, вернут ошибку. В листинге 21-24 мы изменяем цикл Worker, чтобы в этом случае корректно выйти из цикла; это означает, что потоки завершатся, когда реализация drop для ThreadPool вызовет для них join.

Имя файла: src/lib.rs
use std::{
    sync::{Arc, Mutex, mpsc},
    thread,
};

pub struct ThreadPool {
    workers: Vec<Worker>,
    sender: Option<mpsc::Sender<Job>>,
}

type Job = Box<dyn FnOnce() + Send + 'static>;

impl ThreadPool {
    /// Create a new ThreadPool.
    ///
    /// The size is the number of threads in the pool.
    ///
    /// # Panics
    ///
    /// The `new` function will panic if the size is zero.
    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: Some(sender),
        }
    }

    pub fn execute<F>(&self, f: F)
    where
        F: FnOnce() + Send + 'static,
    {
        let job = Box::new(f);

        self.sender.as_ref().unwrap().send(job).unwrap();
    }
}

impl Drop for ThreadPool {
    fn drop(&mut self) {
        drop(self.sender.take());

        for worker in self.workers.drain(..) {
            println!("Shutting down worker {}", worker.id);

            worker.thread.join().unwrap();
        }
    }
}

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 message = receiver.lock().unwrap().recv();

                match message {
                    Ok(job) => {
                        println!("Worker {id} got a job; executing.");

                        job();
                    }
                    Err(_) => {
                        println!("Worker {id} disconnected; shutting down.");
                        break;
                    }
                }
            }
        });

        Worker { id, thread }
    }
}
Listing 21-24: Явный выход из цикла, когда recv возвращает ошибку

Чтобы увидеть этот код в действии, изменим main, чтобы он принимал только два запроса перед корректным завершением сервера, как показано в листинге 21-25.

Имя файла: src/main.rs
use hello::ThreadPool;
use std::{
    fs,
    io::{BufReader, prelude::*},
    net::{TcpListener, TcpStream},
    thread,
    time::Duration,
};

fn main() {
    let listener = TcpListener::bind("127.0.0.1:7878").unwrap();
    let pool = ThreadPool::new(4);

    for stream in listener.incoming().take(2) {
        let stream = stream.unwrap();

        pool.execute(|| {
            handle_connection(stream);
        });
    }

    println!("Shutting down.");
}

fn handle_connection(mut stream: TcpStream) {
    let buf_reader = BufReader::new(&stream);
    let request_line = buf_reader.lines().next().unwrap().unwrap();

    let (status_line, filename) = match &request_line[..] {
        "GET / HTTP/1.1" => ("HTTP/1.1 200 OK", "hello.html"),
        "GET /sleep HTTP/1.1" => {
            thread::sleep(Duration::from_secs(5));
            ("HTTP/1.1 200 OK", "hello.html")
        }
        _ => ("HTTP/1.1 404 NOT FOUND", "404.html"),
    };

    let contents = fs::read_to_string(filename).unwrap();
    let length = contents.len();

    let response =
        format!("{status_line}\r\nContent-Length: {length}\r\n\r\n{contents}");

    stream.write_all(response.as_bytes()).unwrap();
}
Listing 21-25: Завершение сервера после двух запросов путем выхода из цикла

Вы не хотели бы, чтобы настоящий веб-сервер завершал работу после обслуживания только двух запросов. Этот код лишь демонстрирует, что корректное завершение и очистка работают как положено.

Метод take определен в трейте Iterator и ограничивает итерацию максимум первыми двумя элементами. ThreadPool выйдет из области видимости в конце main, и реализация drop будет выполнена.

Запустите сервер с помощью cargo run и сделайте три запроса. Третий запрос должен завершиться ошибкой, а в терминале вы должны увидеть вывод, похожий на этот:

$ cargo run
   Compiling hello v0.1.0 (file:///projects/hello)
    Finished `dev` profile [unoptimized + debuginfo] target(s) in 0.41s
     Running `target/debug/hello`
Worker 0 got a job; executing.
Shutting down.
Shutting down worker 0
Worker 3 got a job; executing.
Worker 1 disconnected; shutting down.
Worker 2 disconnected; shutting down.
Worker 3 disconnected; shutting down.
Worker 0 disconnected; shutting down.
Shutting down worker 1
Shutting down worker 2
Shutting down worker 3

Вы можете увидеть другой порядок идентификаторов Worker и напечатанных сообщений. Из сообщений видно, как работает этот код: экземпляры Worker 0 и 3 получили первые два запроса. Сервер перестал принимать соединения после второго соединения, и реализация Drop для ThreadPool начинает выполняться еще до того, как Worker 3 начинает свое задание. Сброс sender отключает все экземпляры Worker и сообщает им, что нужно завершить работу. Каждый экземпляр Worker печатает сообщение при отключении, а затем пул потоков вызывает join, чтобы дождаться завершения каждого потока Worker.

Обратите внимание на один интересный аспект именно этого выполнения: ThreadPool сбросил sender, и до того как какой-либо Worker получил ошибку, мы попытались присоединить Worker 0. Worker 0 еще не получил ошибку от recv, поэтому основной поток заблокировался, ожидая завершения Worker 0. Тем временем Worker 3 получил задание, а затем все потоки получили ошибку. Когда Worker 0 завершился, основной поток дождался завершения остальных экземпляров Worker. К этому моменту все они вышли из своих циклов и остановились.

Поздравляем! Теперь мы завершили наш проект: у нас есть базовый веб-сервер, который использует пул потоков, чтобы отвечать асинхронно. Мы можем выполнять корректное завершение сервера, которое очищает все потоки в пуле.

Вот полный код для справки:

Имя файла: src/main.rs
use hello::ThreadPool;
use std::{
    fs,
    io::{BufReader, prelude::*},
    net::{TcpListener, TcpStream},
    thread,
    time::Duration,
};

fn main() {
    let listener = TcpListener::bind("127.0.0.1:7878").unwrap();
    let pool = ThreadPool::new(4);

    for stream in listener.incoming().take(2) {
        let stream = stream.unwrap();

        pool.execute(|| {
            handle_connection(stream);
        });
    }

    println!("Shutting down.");
}

fn handle_connection(mut stream: TcpStream) {
    let buf_reader = BufReader::new(&stream);
    let request_line = buf_reader.lines().next().unwrap().unwrap();

    let (status_line, filename) = match &request_line[..] {
        "GET / HTTP/1.1" => ("HTTP/1.1 200 OK", "hello.html"),
        "GET /sleep HTTP/1.1" => {
            thread::sleep(Duration::from_secs(5));
            ("HTTP/1.1 200 OK", "hello.html")
        }
        _ => ("HTTP/1.1 404 NOT FOUND", "404.html"),
    };

    let contents = fs::read_to_string(filename).unwrap();
    let length = contents.len();

    let response =
        format!("{status_line}\r\nContent-Length: {length}\r\n\r\n{contents}");

    stream.write_all(response.as_bytes()).unwrap();
}
Имя файла: src/lib.rs
use std::{
    sync::{Arc, Mutex, mpsc},
    thread,
};

pub struct ThreadPool {
    workers: Vec<Worker>,
    sender: Option<mpsc::Sender<Job>>,
}

type Job = Box<dyn FnOnce() + Send + 'static>;

impl ThreadPool {
    /// Create a new ThreadPool.
    ///
    /// The size is the number of threads in the pool.
    ///
    /// # Panics
    ///
    /// The `new` function will panic if the size is zero.
    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: Some(sender),
        }
    }

    pub fn execute<F>(&self, f: F)
    where
        F: FnOnce() + Send + 'static,
    {
        let job = Box::new(f);

        self.sender.as_ref().unwrap().send(job).unwrap();
    }
}

impl Drop for ThreadPool {
    fn drop(&mut self) {
        drop(self.sender.take());

        for worker in &mut self.workers {
            println!("Shutting down worker {}", worker.id);

            if let Some(thread) = worker.thread.take() {
                thread.join().unwrap();
            }
        }
    }
}

struct Worker {
    id: usize,
    thread: Option<thread::JoinHandle<()>>,
}

impl Worker {
    fn new(id: usize, receiver: Arc<Mutex<mpsc::Receiver<Job>>>) -> Worker {
        let thread = thread::spawn(move || {
            loop {
                let message = receiver.lock().unwrap().recv();

                match message {
                    Ok(job) => {
                        println!("Worker {id} got a job; executing.");

                        job();
                    }
                    Err(_) => {
                        println!("Worker {id} disconnected; shutting down.");
                        break;
                    }
                }
            }
        });

        Worker {
            id,
            thread: Some(thread),
        }
    }
}

Здесь можно сделать больше! Если вы хотите продолжить улучшать этот проект, вот несколько идей:

  • Добавьте больше документации к ThreadPool и его публичным методам.
  • Добавьте тесты функциональности библиотеки.
  • Замените вызовы unwrap более надежной обработкой ошибок.
  • Используйте ThreadPool для выполнения какой-нибудь задачи, кроме обслуживания веб-запросов.
  • Найдите крейт пула потоков на crates.io и реализуйте похожий веб-сервер с использованием этого крейта. Затем сравните его API и надежность с пулом потоков, который реализовали мы.

Итоги

Отличная работа! Вы добрались до конца книги! Мы хотим поблагодарить вас за то, что присоединились к нам в этом путешествии по Rust. Теперь вы готовы реализовывать собственные проекты на Rust и помогать с проектами других людей. Помните, что существует дружелюбное сообщество других Rustaceans, которые будут рады помочь вам с любыми трудностями, с которыми вы столкнетесь на своем пути в Rust.