Корректное завершение и очистка
Код в листинге 21-20 отвечает на запросы асинхронно с помощью пула потоков, как
мы и планировали. Мы получаем несколько предупреждений о полях workers, id
и thread, которые не используем напрямую; они напоминают нам, что мы ничего
не очищаем. Когда мы используем менее изящный метод
ctrl-C, чтобы остановить основной поток, все остальные
потоки тоже останавливаются немедленно, даже если они находятся в середине
обслуживания запроса.
Далее мы реализуем трейт Drop, чтобы вызывать join для каждого потока в
пуле, позволяя им завершить запросы, над которыми они работают, прежде чем
закрыться. Затем мы реализуем способ сообщать потокам, что им нужно перестать
принимать новые запросы и завершить работу. Чтобы увидеть этот код в действии,
мы изменим сервер так, чтобы он принимал только два запроса перед корректным
завершением своего пула потоков.
Стоит заметить по ходу дела: ничто из этого не влияет на части кода, которые отвечают за выполнение замыканий, поэтому все здесь было бы таким же, если бы мы использовали пул потоков для async runtime.
Реализация трейта Drop для ThreadPool
Начнем с реализации Drop для нашего пула потоков. Когда пул сбрасывается, все
наши потоки должны присоединиться, чтобы гарантировать завершение своей работы.
Листинг 21-22 показывает первую попытку реализации Drop; этот код пока не
совсем будет работать.
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 }
}
}
Сначала мы проходим циклом по всем 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 вот так:
#![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.
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 }
}
}
sender перед присоединением потоков WorkerСброс sender закрывает канал, что указывает: больше сообщений отправляться не
будет. Когда это происходит, все вызовы recv, которые экземпляры Worker
делают в бесконечном цикле, вернут ошибку. В листинге 21-24 мы изменяем цикл
Worker, чтобы в этом случае корректно выйти из цикла; это означает, что
потоки завершатся, когда реализация drop для ThreadPool вызовет для них
join.
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 }
}
}
recv возвращает ошибкуЧтобы увидеть этот код в действии, изменим main, чтобы он принимал только два
запроса перед корректным завершением сервера, как показано в листинге 21-25.
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();
}
Вы не хотели бы, чтобы настоящий веб-сервер завершал работу после обслуживания только двух запросов. Этот код лишь демонстрирует, что корректное завершение и очистка работают как положено.
Метод 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. К этому моменту все они вышли из своих циклов и
остановились.
Поздравляем! Теперь мы завершили наш проект: у нас есть базовый веб-сервер, который использует пул потоков, чтобы отвечать асинхронно. Мы можем выполнять корректное завершение сервера, которое очищает все потоки в пуле.
Вот полный код для справки:
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();
}
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.