От однопоточного сервера к многопоточному
Сейчас сервер будет обрабатывать каждый запрос по очереди, то есть не начнет обрабатывать второе соединение, пока не закончит обработку первого. Если сервер будет получать все больше и больше запросов, такое последовательное выполнение будет становиться все менее оптимальным. Если сервер получит запрос, обработка которого занимает много времени, последующим запросам придется ждать, пока долгий запрос не завершится, даже если новые запросы можно обработать быстро. Нам нужно это исправить, но сначала посмотрим на проблему в действии.
Имитация медленного запроса
Мы посмотрим, как медленно обрабатываемый запрос может повлиять на другие запросы к текущей реализации нашего сервера. Листинг 21-10 реализует обработку запроса к /sleep с имитированным медленным ответом, из-за которого сервер заснет на пять секунд перед ответом.
use std::{
fs,
io::{BufReader, prelude::*},
net::{TcpListener, TcpStream},
thread,
time::Duration,
};
// --snip--
fn main() {
let listener = TcpListener::bind("127.0.0.1:7878").unwrap();
for stream in listener.incoming() {
let stream = stream.unwrap();
handle_connection(stream);
}
}
fn handle_connection(mut stream: TcpStream) {
// --snip--
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"),
};
// --snip--
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();
}
Теперь мы перешли с if на match, потому что у нас есть три случая. Нам
нужно явно сопоставлять срез request_line со строковыми литералами; match
не выполняет автоматическое взятие ссылки и разыменование так, как это делает
метод проверки равенства.
Первая ветвь такая же, как блок if из листинга 21-9. Вторая ветвь
соответствует запросу к /sleep. Когда такой запрос получен, сервер засыпает
на пять секунд перед отрисовкой успешной HTML-страницы. Третья ветвь такая же,
как блок else из листинга 21-9.
Вы можете видеть, насколько примитивен наш сервер: настоящие библиотеки обрабатывали бы распознавание нескольких запросов гораздо менее многословно!
Запустите сервер с помощью cargo run. Затем откройте два окна браузера: одно
для http://127.0.0.1:7878, а другое для
http://127.0.0.1:7878/sleep. Если вы несколько раз введете URI /, как
раньше, вы увидите, что он отвечает быстро. Но если вы введете /sleep, а
затем загрузите /, то увидите, что / ждет, пока sleep полностью проспит
свои пять секунд, прежде чем загрузиться.
Есть несколько техник, которые мы могли бы использовать, чтобы избежать накопления запросов за медленным запросом, включая использование async, как мы делали в главе 17; та, которую мы реализуем, — это пул потоков.
Улучшение пропускной способности с помощью пула потоков
Пул потоков — это группа созданных потоков, которые готовы и ждут обработки задачи. Когда программа получает новую задачу, она назначает один из потоков в пуле этой задаче, и этот поток обрабатывает задачу. Остальные потоки в пуле доступны для обработки любых других задач, которые поступят, пока первый поток занят обработкой. Когда первый поток заканчивает обработку своей задачи, он возвращается в пул простаивающих потоков, готовый обработать новую задачу. Пул потоков позволяет обрабатывать соединения конкурентно, увеличивая пропускную способность вашего сервера.
Мы ограничим количество потоков в пуле небольшим числом, чтобы защититься от DoS-атак; если бы наша программа создавала новый поток для каждого входящего запроса, кто-то, отправивший 10 миллионов запросов к нашему серверу, мог бы устроить хаос, израсходовав все ресурсы сервера и полностью остановив обработку запросов.
Итак, вместо создания неограниченного числа потоков у нас будет фиксированное
число потоков, ожидающих в пуле. Входящие запросы отправляются в пул для
обработки. Пул будет поддерживать очередь входящих запросов. Каждый поток в
пуле будет извлекать запрос из этой очереди, обрабатывать запрос, а затем
просить у очереди другой запрос. С такой схемой мы можем обрабатывать до
N запросов конкурентно, где N — количество потоков. Если каждый поток
отвечает на долго выполняющийся запрос, последующие запросы все равно могут
накапливаться в очереди, но мы увеличили количество долгих запросов, которые
можем обработать до достижения этой точки.
Эта техника — лишь один из многих способов улучшить пропускную способность веб-сервера. Другие варианты, которые можно изучить, — модель fork/join, однопоточная модель асинхронного ввода-вывода и многопоточная модель асинхронного ввода-вывода. Если вам интересна эта тема, можно подробнее почитать о других решениях и попробовать реализовать их; с низкоуровневым языком вроде Rust возможны все эти варианты.
Прежде чем начать реализовывать пул потоков, поговорим о том, как должно выглядеть использование пула. Когда вы пытаетесь проектировать код, написание клиентского интерфейса первым может помочь направить дизайн. Напишите API кода так, чтобы он был структурирован так, как вы хотите его вызывать; затем реализуйте функциональность внутри этой структуры, а не реализуйте функциональность и только потом проектируйте публичный API.
Подобно тому как в проекте главы 12 мы использовали разработку через тестирование, здесь мы будем использовать разработку, направляемую компилятором. Мы напишем код, который вызывает нужные нам функции, а затем посмотрим на ошибки компилятора, чтобы определить, что нужно изменить дальше, чтобы код заработал. Однако прежде чем это сделать, мы изучим технику, которую не собираемся использовать как отправную точку.
Создание потока для каждого запроса
Сначала рассмотрим, как мог бы выглядеть наш код, если бы он создавал новый поток для каждого соединения. Как упоминалось ранее, это не наш окончательный план из-за проблем с потенциальным созданием неограниченного количества потоков, но это отправная точка, чтобы сначала получить работающий многопоточный сервер. Затем мы добавим пул потоков как улучшение, и сравнивать два решения будет проще.
Листинг 21-11 показывает изменения, которые нужно внести в main, чтобы
создавать новый поток для обработки каждого потока данных внутри цикла for.
use std::{
fs,
io::{BufReader, prelude::*},
net::{TcpListener, TcpStream},
thread,
time::Duration,
};
fn main() {
let listener = TcpListener::bind("127.0.0.1:7878").unwrap();
for stream in listener.incoming() {
let stream = stream.unwrap();
thread::spawn(|| {
handle_connection(stream);
});
}
}
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();
}
Как вы узнали в главе 16, thread::spawn создаст новый поток, а затем
запустит код в замыкании в новом потоке. Если вы запустите этот код и загрузите
/sleep в браузере, а затем / еще в двух вкладках браузера, вы действительно
увидите, что запросам к / не нужно ждать завершения /sleep. Однако, как мы
уже упоминали, со временем это перегрузит систему, потому что вы будете
создавать новые потоки без какого-либо ограничения.
Возможно, вы также помните из главы 17, что это именно такая ситуация, где async и await действительно сияют! Держите это в уме, пока мы строим пул потоков, и подумайте, как все выглядело бы иначе или так же с async.
Создание конечного числа потоков
Мы хотим, чтобы наш пул потоков работал похожим и привычным образом, чтобы
переход от потоков к пулу потоков не требовал больших изменений в коде,
который использует наш API. Листинг 21-12 показывает гипотетический интерфейс
для структуры ThreadPool, которую мы хотим использовать вместо
thread::spawn.
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() {
let stream = stream.unwrap();
pool.execute(|| {
handle_connection(stream);
});
}
}
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();
}
ThreadPoolМы используем ThreadPool::new, чтобы создать новый пул потоков с настраиваемым
количеством потоков, в данном случае четырьмя. Затем в цикле for
pool.execute имеет интерфейс, похожий на thread::spawn, поскольку принимает
замыкание, которое пул должен выполнить для каждого потока данных. Нам нужно
реализовать pool.execute так, чтобы он принимал замыкание и передавал его
потоку в пуле для выполнения. Этот код пока не скомпилируется, но мы попробуем,
чтобы компилятор мог подсказать, как это исправить.
Создание ThreadPool с помощью разработки, направляемой компилятором
Внесите изменения из листинга 21-12 в src/main.rs, а затем используем ошибки
компилятора из cargo check, чтобы направлять нашу разработку. Вот первая
ошибка, которую мы получим:
$ cargo check
Checking hello v0.1.0 (file:///projects/hello)
error[E0433]: failed to resolve: use of undeclared type `ThreadPool`
--> src/main.rs:11:16
|
11 | let pool = ThreadPool::new(4);
| ^^^^^^^^^^ use of undeclared type `ThreadPool`
For more information about this error, try `rustc --explain E0433`.
error: could not compile `hello` (bin "hello") due to 1 previous error
Отлично! Эта ошибка говорит нам, что нужен тип или модуль ThreadPool, поэтому
сейчас мы его создадим. Наша реализация ThreadPool будет независима от вида
работы, которую выполняет веб-сервер. Поэтому переключим крейт hello с
бинарного крейта на библиотечный крейт, чтобы хранить в нем реализацию
ThreadPool. После перехода к библиотечному крейту мы также могли бы
использовать отдельную библиотеку пула потоков для любой работы, которую хотим
выполнять с помощью пула потоков, а не только для обслуживания веб-запросов.
Создайте файл src/lib.rs, содержащий следующее, что пока является самым
простым определением структуры ThreadPool, какое у нас может быть:
pub struct ThreadPool;
Затем отредактируйте файл main.rs, чтобы ввести ThreadPool из библиотечного
крейта в область видимости, добавив следующий код в начало 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() {
let stream = stream.unwrap();
pool.execute(|| {
handle_connection(stream);
});
}
}
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();
}
Этот код все еще не будет работать, но проверим его снова, чтобы получить следующую ошибку, которую нужно исправить:
$ cargo check
Checking hello v0.1.0 (file:///projects/hello)
error[E0599]: no function or associated item named `new` found for struct `ThreadPool` in the current scope
--> src/main.rs:12:28
|
12 | let pool = ThreadPool::new(4);
| ^^^ function or associated item not found in `ThreadPool`
For more information about this error, try `rustc --explain E0599`.
error: could not compile `hello` (bin "hello") due to 1 previous error
Эта ошибка указывает, что дальше нам нужно создать связанную функцию с именем
new для ThreadPool. Мы также знаем, что new должна иметь один параметр,
который может принять 4 как аргумент, и должна возвращать экземпляр
ThreadPool. Реализуем самую простую функцию new, которая будет обладать
этими характеристиками:
pub struct ThreadPool;
impl ThreadPool {
pub fn new(size: usize) -> ThreadPool {
ThreadPool
}
}
Мы выбрали usize как тип параметра size, потому что знаем, что отрицательное
число потоков не имеет смысла. Мы также знаем, что будем использовать это 4
как количество элементов в коллекции потоков, а именно для этого и предназначен
тип usize, как обсуждалось в разделе «Целочисленные
типы» главы 3.
Проверим код снова:
$ cargo check
Checking hello v0.1.0 (file:///projects/hello)
error[E0599]: no method named `execute` found for struct `ThreadPool` in the current scope
--> src/main.rs:17:14
|
17 | pool.execute(|| {
| -----^^^^^^^ method not found in `ThreadPool`
For more information about this error, try `rustc --explain E0599`.
error: could not compile `hello` (bin "hello") due to 1 previous error
Теперь ошибка возникает потому, что у нас нет метода execute у ThreadPool.
Вспомните из раздела «Создание конечного числа
потоков», что мы решили:
наш пул потоков должен иметь интерфейс, похожий на thread::spawn. Кроме
того, мы реализуем функцию execute так, чтобы она принимала переданное ей
замыкание и отдавала его свободному потоку в пуле для выполнения.
Мы определим метод execute у ThreadPool, который принимает замыкание как
параметр. Вспомните из раздела «Перемещение захваченных значений из
замыканий» главы 13, что мы можем
принимать замыкания как параметры с тремя разными трейтами: Fn, FnMut и
FnOnce. Нам нужно решить, какой вид замыкания использовать здесь. Мы знаем,
что в итоге будем делать что-то похожее на реализацию thread::spawn из
стандартной библиотеки, поэтому можем посмотреть, какие ограничения есть у
сигнатуры thread::spawn на ее параметр. Документация показывает следующее:
pub fn spawn<F, T>(f: F) -> JoinHandle<T>
where
F: FnOnce() -> T,
F: Send + 'static,
T: Send + 'static,
Параметр типа F — тот, который нас здесь интересует; параметр типа T
связан с возвращаемым значением, и сейчас он нас не волнует. Мы видим, что
spawn использует FnOnce как ограничение трейта для F. Вероятно, это то,
что нужно и нам, потому что в итоге мы передадим аргумент, полученный в
execute, в spawn. Мы можем быть еще более уверены, что FnOnce — нужный
нам трейт, потому что поток для выполнения запроса выполнит замыкание этого
запроса только один раз, что соответствует Once в FnOnce.
У параметра типа F также есть ограничение трейта Send и ограничение времени
жизни 'static, которые полезны в нашей ситуации: Send нужен, чтобы
передать замыкание из одного потока в другой, а 'static — потому что мы не
знаем, сколько времени поток будет выполнять замыкание. Создадим метод
execute у ThreadPool, который будет принимать обобщенный параметр типа F
с этими ограничениями:
pub struct ThreadPool;
impl ThreadPool {
// --snip--
pub fn new(size: usize) -> ThreadPool {
ThreadPool
}
pub fn execute<F>(&self, f: F)
where
F: FnOnce() + Send + 'static,
{
}
}
Мы все еще используем () после FnOnce, потому что этот FnOnce
представляет замыкание, которое не принимает параметров и возвращает unit-тип
(). Как и в определениях функций, возвращаемый тип можно опустить из
сигнатуры, но даже если параметров нет, круглые скобки все равно нужны.
Опять же, это самая простая реализация метода execute: она ничего не делает,
но мы только пытаемся добиться компиляции кода. Проверим его снова:
$ cargo check
Checking hello v0.1.0 (file:///projects/hello)
Finished `dev` profile [unoptimized + debuginfo] target(s) in 0.24s
Он компилируется! Но обратите внимание: если вы попробуете cargo run и
сделаете запрос в браузере, то увидите в браузере ошибки, которые мы видели в
начале главы. Наша библиотека пока на самом деле не вызывает замыкание,
переданное в execute!
Примечание: поговорка, которую можно услышать о языках со строгими компиляторами, таких как Haskell и Rust, звучит так: «Если код компилируется, он работает». Но эта поговорка не является универсально верной. Наш проект компилируется, но он не делает абсолютно ничего! Если бы мы строили настоящий, полный проект, сейчас было бы хорошее время начать писать модульные тесты, чтобы проверить, что код компилируется и ведет себя так, как мы хотим.
Подумайте: что было бы здесь иначе, если бы мы собирались выполнять future вместо замыкания?
Проверка количества потоков в new
Мы ничего не делаем с параметрами new и execute. Давайте реализуем тела
этих функций с нужным нам поведением. Для начала подумаем о new. Ранее мы
выбрали беззнаковый тип для параметра size, потому что пул с отрицательным
числом потоков не имеет смысла. Однако пул с нулем потоков тоже не имеет
смысла, но ноль является вполне допустимым значением usize. Мы добавим код,
проверяющий, что size больше нуля перед возвратом экземпляра ThreadPool, и
заставим программу паниковать, если она получит ноль, используя макрос
assert!, как показано в листинге 21-13.
pub struct ThreadPool;
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);
ThreadPool
}
// --snip--
pub fn execute<F>(&self, f: F)
where
F: FnOnce() + Send + 'static,
{
}
}
ThreadPool::new, вызывающая панику при нулевом sizeМы также добавили документацию для нашего ThreadPool с помощью doc-комментариев.
Обратите внимание, что мы следовали хорошим практикам документирования,
добавив раздел, в котором указаны ситуации, когда функция может паниковать,
как обсуждалось в главе 14. Попробуйте выполнить cargo doc --open и
щелкнуть структуру ThreadPool, чтобы увидеть, как выглядит сгенерированная
документация для new!
Вместо добавления макроса assert!, как мы сделали здесь, мы могли бы
изменить new на build и возвращать Result, как делали с Config::build
в проекте ввода-вывода в листинге 12-9. Но в этом случае мы решили, что
попытка создать пул потоков без потоков должна быть невосстановимой ошибкой.
Если вы чувствуете в себе силы, попробуйте написать функцию с именем build
со следующей сигнатурой, чтобы сравнить ее с функцией new:
pub fn build(size: usize) -> Result<ThreadPool, PoolCreationError> {
Создание места для хранения потоков
Теперь, когда у нас есть способ убедиться, что у нас допустимое количество
потоков для хранения в пуле, мы можем создать эти потоки и сохранить их в
структуре ThreadPool перед возвратом структуры. Но как «хранить» поток?
Давайте еще раз посмотрим на сигнатуру thread::spawn:
pub fn spawn<F, T>(f: F) -> JoinHandle<T>
where
F: FnOnce() -> T,
F: Send + 'static,
T: Send + 'static,
Функция spawn возвращает JoinHandle<T>, где T — тип, возвращаемый
замыканием. Попробуем тоже использовать JoinHandle и посмотрим, что
произойдет. В нашем случае замыкания, которые мы передаем пулу потоков, будут
обрабатывать соединение и ничего не возвращать, поэтому T будет unit-типом
().
Код в листинге 21-14 скомпилируется, но пока не создает никаких потоков. Мы
изменили определение ThreadPool, чтобы оно содержало вектор экземпляров
thread::JoinHandle<()>, инициализировали вектор с емкостью size, настроили
цикл for, который будет выполнять некоторый код для создания потоков, и
вернули экземпляр ThreadPool, содержащий их.
use std::thread;
pub struct ThreadPool {
threads: Vec<thread::JoinHandle<()>>,
}
impl ThreadPool {
// --snip--
/// 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 mut threads = Vec::with_capacity(size);
for _ in 0..size {
// create some threads and store them in the vector
}
ThreadPool { threads }
}
// --snip--
pub fn execute<F>(&self, f: F)
where
F: FnOnce() + Send + 'static,
{
}
}
ThreadPoolМы ввели std::thread в область видимости в библиотечном крейте, потому что
используем thread::JoinHandle как тип элементов вектора в ThreadPool.
После получения допустимого размера наш ThreadPool создает новый вектор,
который может содержать size элементов. Функция with_capacity выполняет ту
же задачу, что и Vec::new, но с важным отличием: она заранее выделяет место в
векторе. Поскольку мы знаем, что нужно хранить size элементов в векторе,
такое выделение памяти заранее немного эффективнее, чем использование
Vec::new, который изменяет свой размер по мере вставки элементов.
Когда вы снова выполните cargo check, он должен завершиться успешно.
Отправка кода из ThreadPool в поток
В цикле for в листинге 21-14 мы оставили комментарий о создании потоков.
Здесь мы посмотрим, как на самом деле создавать потоки. Стандартная библиотека
предоставляет thread::spawn как способ создания потоков, и thread::spawn
ожидает получить код, который поток должен выполнить сразу после создания.
Однако в нашем случае мы хотим создать потоки и заставить их ждать код,
который мы отправим позже. Реализация потоков в стандартной библиотеке не
включает способ сделать это; нам придется реализовать его вручную.
Мы реализуем это поведение, введя новую структуру данных между ThreadPool и
потоками, которая будет управлять этим новым поведением. Мы назовем эту
структуру данных Worker, что является распространенным термином в реализациях
пулов. Worker забирает код, который нужно выполнить, и выполняет этот код в
своем потоке.
Представьте людей, работающих на кухне в ресторане: работники ждут, пока поступят заказы от клиентов, а затем отвечают за то, чтобы взять эти заказы и выполнить их.
Вместо хранения вектора экземпляров JoinHandle<()> в пуле потоков мы будем
хранить экземпляры структуры Worker. Каждый Worker будет хранить один
экземпляр JoinHandle<()>. Затем мы реализуем метод у Worker, который будет
принимать замыкание с кодом для выполнения и отправлять его уже запущенному
потоку для выполнения. Мы также дадим каждому Worker id, чтобы при
журналировании или отладке различать разные экземпляры Worker в пуле.
Вот новый процесс, который будет происходить при создании ThreadPool. Код,
отправляющий замыкание в поток, мы реализуем после того, как настроим Worker
таким образом:
- Определить структуру
Worker, которая хранитidиJoinHandle<()>. - Изменить
ThreadPool, чтобы он хранил вектор экземпляровWorker. - Определить функцию
Worker::new, которая принимает числоidи возвращает экземплярWorker, хранящийidи поток, созданный с пустым замыканием. - В
ThreadPool::newиспользовать счетчик циклаfor, чтобы сгенерироватьid, создать новыйWorkerс этимidи сохранитьWorkerв векторе.
Если вы готовы к испытанию, попробуйте реализовать эти изменения самостоятельно перед тем, как смотреть код в листинге 21-15.
Готовы? Вот листинг 21-15 с одним из способов внести предыдущие изменения.
use std::thread;
pub struct ThreadPool {
workers: Vec<Worker>,
}
impl ThreadPool {
// --snip--
/// 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 mut workers = Vec::with_capacity(size);
for id in 0..size {
workers.push(Worker::new(id));
}
ThreadPool { workers }
}
// --snip--
pub fn execute<F>(&self, f: F)
where
F: FnOnce() + Send + 'static,
{
}
}
struct Worker {
id: usize,
thread: thread::JoinHandle<()>,
}
impl Worker {
fn new(id: usize) -> Worker {
let thread = thread::spawn(|| {});
Worker { id, thread }
}
}
ThreadPool для хранения экземпляров Worker вместо потоков напрямуюМы изменили имя поля в ThreadPool с threads на workers, потому что теперь
оно хранит экземпляры Worker, а не экземпляры JoinHandle<()>. Мы используем
счетчик в цикле for как аргумент для Worker::new и сохраняем каждый новый
Worker в векторе с именем workers.
Внешнему коду (например, нашему серверу в src/main.rs) не нужно знать детали
реализации, связанные с использованием структуры Worker внутри ThreadPool,
поэтому мы делаем структуру Worker и ее функцию new приватными. Функция
Worker::new использует переданный ей id и сохраняет экземпляр
JoinHandle<()>, созданный путем запуска нового потока с пустым замыканием.
Примечание: если операционная система не может создать поток из-за нехватки системных ресурсов,
thread::spawnвызовет панику. Это приведет к панике всего нашего сервера, даже если создание некоторых потоков могло успешно завершиться. Ради простоты такое поведение приемлемо, но в промышленной реализации пула потоков вы, скорее всего, захотели бы использоватьstd::thread::Builderи его методspawn, который возвращаетResult.
Этот код скомпилируется и сохранит то количество экземпляров Worker, которое
мы указали как аргумент для ThreadPool::new. Но мы все еще не обрабатываем
замыкание, которое получаем в execute. Далее посмотрим, как это сделать.
Отправка запросов потокам через каналы
Следующая проблема, которую мы решим, состоит в том, что замыкания, переданные
в thread::spawn, не делают абсолютно ничего. Сейчас мы получаем замыкание,
которое хотим выполнить, в методе execute. Но нам нужно передать
thread::spawn замыкание для запуска, когда мы создаем каждый Worker во
время создания ThreadPool.
Мы хотим, чтобы только что созданные структуры Worker получали код для
выполнения из очереди, хранящейся в ThreadPool, и отправляли этот код своему
потоку для выполнения.
Каналы, о которых мы узнали в главе 16, — простой способ взаимодействия между
двумя потоками — идеально подходят для этого случая. Мы будем использовать
канал как очередь заданий, а execute будет отправлять задание из ThreadPool
экземплярам Worker, которые отправят задание своему потоку. Вот план:
ThreadPoolсоздаст канал и сохранит отправителя.- Каждый
Workerсохранит получателя. - Мы создадим новую структуру
Job, которая будет хранить замыкания, которые мы хотим отправлять по каналу. - Метод
executeотправит задание, которое хочет выполнить, через отправителя. - В своем потоке
Workerбудет циклически получать данные от своего получателя и выполнять замыкания любых полученных заданий.
Начнем с создания канала в ThreadPool::new и хранения отправителя в
экземпляре ThreadPool, как показано в листинге 21-16. Структура Job пока
ничего не хранит, но будет типом элементов, которые мы отправляем по каналу.
use std::{sync::mpsc, thread};
pub struct ThreadPool {
workers: Vec<Worker>,
sender: mpsc::Sender<Job>,
}
struct Job;
impl ThreadPool {
// --snip--
/// 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 mut workers = Vec::with_capacity(size);
for id in 0..size {
workers.push(Worker::new(id));
}
ThreadPool { workers, sender }
}
// --snip--
pub fn execute<F>(&self, f: F)
where
F: FnOnce() + Send + 'static,
{
}
}
struct Worker {
id: usize,
thread: thread::JoinHandle<()>,
}
impl Worker {
fn new(id: usize) -> Worker {
let thread = thread::spawn(|| {});
Worker { id, thread }
}
}
ThreadPool отправителя канала для экземпляров JobВ ThreadPool::new мы создаем новый канал и заставляем пул хранить
отправителя. Это успешно скомпилируется.
Попробуем передать получателя канала каждому Worker во время создания канала
пулом потоков. Мы знаем, что хотим использовать получателя в потоке, который
запускают экземпляры Worker, поэтому сошлемся на параметр receiver в
замыкании. Код в листинге 21-17 пока не совсем скомпилируется.
use std::{sync::mpsc, thread};
pub struct ThreadPool {
workers: Vec<Worker>,
sender: mpsc::Sender<Job>,
}
struct Job;
impl ThreadPool {
// --snip--
/// 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 mut workers = Vec::with_capacity(size);
for id in 0..size {
workers.push(Worker::new(id, receiver));
}
ThreadPool { workers, sender }
}
// --snip--
pub fn execute<F>(&self, f: F)
where
F: FnOnce() + Send + 'static,
{
}
}
// --snip--
struct Worker {
id: usize,
thread: thread::JoinHandle<()>,
}
impl Worker {
fn new(id: usize, receiver: mpsc::Receiver<Job>) -> Worker {
let thread = thread::spawn(|| {
receiver;
});
Worker { id, thread }
}
}
WorkerМы внесли несколько небольших и прямолинейных изменений: передаем получателя в
Worker::new, а затем используем его внутри замыкания.
Когда мы пытаемся проверить этот код, получаем такую ошибку:
$ cargo check
Checking hello v0.1.0 (file:///projects/hello)
error[E0382]: use of moved value: `receiver`
--> src/lib.rs:26:42
|
21 | let (sender, receiver) = mpsc::channel();
| -------- move occurs because `receiver` has type `std::sync::mpsc::Receiver<Job>`, which does not implement the `Copy` trait
...
25 | for id in 0..size {
| ----------------- inside of this loop
26 | workers.push(Worker::new(id, receiver));
| ^^^^^^^^ value moved here, in previous iteration of loop
|
note: consider changing this parameter type in method `new` to borrow instead if owning the value isn't necessary
--> src/lib.rs:47:33
|
47 | fn new(id: usize, receiver: mpsc::Receiver<Job>) -> Worker {
| --- in this method ^^^^^^^^^^^^^^^^^^^ this parameter takes ownership of the value
help: consider moving the expression out of the loop so it is only moved once
|
25 ~ let mut value = Worker::new(id, receiver);
26 ~ for id in 0..size {
27 ~ workers.push(value);
|
For more information about this error, try `rustc --explain E0382`.
error: could not compile `hello` (lib) due to 1 previous error
Код пытается передать receiver нескольким экземплярам Worker. Это не
сработает, как вы помните из главы 16: реализация каналов, которую
предоставляет Rust, — это несколько производителей и один потребитель. Это
означает, что мы не можем просто клонировать принимающий конец канала, чтобы
исправить этот код. Мы также не хотим отправлять сообщение несколько раз
нескольким потребителям; мы хотим один список сообщений с несколькими
экземплярами Worker, так чтобы каждое сообщение обрабатывалось один раз.
Кроме того, снятие задания с очереди канала включает изменение receiver,
поэтому потокам нужен безопасный способ совместно использовать и изменять
receiver; иначе мы могли бы получить гонки данных (как обсуждалось в главе
16).
Вспомните потокобезопасные умные указатели, обсуждавшиеся в главе 16: чтобы
разделять владение между несколькими потоками и позволять потокам изменять
значение, нужно использовать Arc<Mutex<T>>. Тип Arc позволит нескольким
экземплярам Worker владеть получателем, а Mutex гарантирует, что только
один Worker за раз получает задание от получателя. Листинг 21-18 показывает
изменения, которые нужно внести.
use std::{
sync::{Arc, Mutex, mpsc},
thread,
};
// --snip--
pub struct ThreadPool {
workers: Vec<Worker>,
sender: mpsc::Sender<Job>,
}
struct Job;
impl ThreadPool {
// --snip--
/// 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 }
}
// --snip--
pub fn execute<F>(&self, f: F)
where
F: FnOnce() + Send + 'static,
{
}
}
// --snip--
struct Worker {
id: usize,
thread: thread::JoinHandle<()>,
}
impl Worker {
fn new(id: usize, receiver: Arc<Mutex<mpsc::Receiver<Job>>>) -> Worker {
// --snip--
let thread = thread::spawn(|| {
receiver;
});
Worker { id, thread }
}
}
Arc и MutexВ ThreadPool::new мы помещаем получателя в Arc и Mutex. Для каждого
нового Worker мы клонируем Arc, чтобы увеличить счетчик ссылок, благодаря
чему экземпляры Worker могут совместно владеть получателем.
С этими изменениями код компилируется! Мы приближаемся!
Реализация метода execute
Наконец реализуем метод execute у ThreadPool. Мы также изменим Job со
структуры на псевдоним типа для трейт-объекта, который хранит тип замыкания,
получаемого execute. Как обсуждалось в разделе «Синонимы типов и
псевдонимы типов» главы 20, псевдонимы типов
позволяют нам сокращать длинные типы для удобства использования. Посмотрите на
листинг 21-19.
use std::{
sync::{Arc, Mutex, mpsc},
thread,
};
pub struct ThreadPool {
workers: Vec<Worker>,
sender: mpsc::Sender<Job>,
}
// --snip--
type Job = Box<dyn FnOnce() + Send + 'static>;
impl ThreadPool {
// --snip--
/// 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();
}
}
// --snip--
struct Worker {
id: usize,
thread: thread::JoinHandle<()>,
}
impl Worker {
fn new(id: usize, receiver: Arc<Mutex<mpsc::Receiver<Job>>>) -> Worker {
let thread = thread::spawn(|| {
receiver;
});
Worker { id, thread }
}
}
Job для Box с замыканием и отправка задания по каналуПосле создания нового экземпляра Job с помощью замыкания, полученного в
execute, мы отправляем это задание через отправляющий конец канала. Мы
вызываем unwrap у send на случай, если отправка завершится неудачно. Это
может произойти, например, если мы остановим выполнение всех наших потоков, то
есть принимающий конец перестанет получать новые сообщения. Сейчас мы не можем
остановить выполнение потоков: наши потоки продолжают выполняться, пока
существует пул. Причина, по которой мы используем unwrap, в том, что мы
знаем: случай ошибки не произойдет, но компилятор этого не знает.
Но мы еще не совсем закончили! В Worker замыкание, передаваемое
thread::spawn, все еще только ссылается на принимающий конец канала.
Вместо этого нам нужно, чтобы замыкание выполняло бесконечный цикл, запрашивая
у принимающего конца канала задание и выполняя его, когда оно получено. Внесем
изменение в Worker::new, показанное в листинге 21-20.
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();
}
}
struct Worker {
id: usize,
thread: thread::JoinHandle<()>,
}
// --snip--
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 }
}
}
WorkerЗдесь мы сначала вызываем lock у receiver, чтобы получить мьютекс, а затем
вызываем unwrap, чтобы паниковать при любых ошибках. Получение блокировки
может завершиться неудачно, если мьютекс находится в отравленном состоянии,
что может случиться, если какой-то другой поток запаниковал, удерживая
блокировку, вместо того чтобы освободить ее. В такой ситуации вызов unwrap,
чтобы этот поток запаниковал, — правильное действие. Можете заменить этот
unwrap на expect с сообщением об ошибке, которое для вас осмысленно.
Если мы получаем блокировку мьютекса, то вызываем recv, чтобы получить Job
из канала. Последний unwrap также пропускает любые ошибки здесь; они могут
возникнуть, если поток, владеющий отправителем, завершил работу, подобно тому,
как метод send возвращает Err, если получатель завершил работу.
Вызов recv блокирует выполнение, поэтому, если задания пока нет, текущий
поток будет ждать, пока задание не станет доступным. Mutex<T> гарантирует,
что только один поток Worker за раз пытается запросить задание.
Наш пул потоков теперь находится в рабочем состоянии! Выполните cargo run и
сделайте несколько запросов:
$ cargo run
Compiling hello v0.1.0 (file:///projects/hello)
warning: field `workers` is never read
--> src/lib.rs:7:5
|
6 | pub struct ThreadPool {
| ---------- field in this struct
7 | workers: Vec<Worker>,
| ^^^^^^^
|
= note: `#[warn(dead_code)]` on by default
warning: fields `id` and `thread` are never read
--> src/lib.rs:48:5
|
47 | struct Worker {
| ------ fields in this struct
48 | id: usize,
| ^^
49 | thread: thread::JoinHandle<()>,
| ^^^^^^
warning: `hello` (lib) generated 2 warnings
Finished `dev` profile [unoptimized + debuginfo] target(s) in 4.91s
Running `target/debug/hello`
Worker 0 got a job; executing.
Worker 2 got a job; executing.
Worker 1 got a job; executing.
Worker 3 got a job; executing.
Worker 0 got a job; executing.
Worker 2 got a job; executing.
Worker 1 got a job; executing.
Worker 3 got a job; executing.
Worker 0 got a job; executing.
Worker 2 got a job; executing.
Успех! Теперь у нас есть пул потоков, который выполняет соединения асинхронно. Никогда не создается больше четырех потоков, поэтому наша система не будет перегружена, если сервер получит много запросов. Если мы сделаем запрос к /sleep, сервер сможет обслуживать другие запросы, выполняя их в другом потоке.
Примечание: если открыть /sleep одновременно в нескольких окнах браузера, они могут загружаться по одному с интервалом в пять секунд. Некоторые веб-браузеры выполняют несколько экземпляров одного и того же запроса последовательно из соображений кеширования. Это ограничение вызвано не нашим веб-сервером.
Сейчас хороший момент остановиться и подумать, чем отличался бы код в листингах 21-18, 21-19 и 21-20, если бы мы использовали futures вместо замыкания для работы, которую нужно выполнить. Какие типы изменились бы? Как изменились бы сигнатуры методов, если вообще изменились бы? Какие части кода остались бы такими же?
После изучения цикла while let в главах 17 и 19 вы, возможно, задаетесь
вопросом, почему мы не написали код потока Worker, как показано в листинге
21-21.
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();
}
}
struct Worker {
id: usize,
thread: thread::JoinHandle<()>,
}
// --snip--
impl Worker {
fn new(id: usize, receiver: Arc<Mutex<mpsc::Receiver<Job>>>) -> Worker {
let thread = thread::spawn(move || {
while let Ok(job) = receiver.lock().unwrap().recv() {
println!("Worker {id} got a job; executing.");
job();
}
});
Worker { id, thread }
}
}
Worker::new с использованием while letЭтот код компилируется и запускается, но не приводит к желаемому поведению
потоков: медленный запрос все равно заставит другие запросы ждать обработки.
Причина довольно тонкая: у структуры Mutex нет публичного метода unlock,
потому что владение блокировкой основано на времени жизни MutexGuard<T>
внутри LockResult<MutexGuard<T>>, который возвращает метод lock. Во время
компиляции проверяющий заимствования затем может обеспечить правило, что
ресурс, защищенный Mutex, нельзя использовать, если мы не удерживаем
блокировку. Однако эта реализация также может привести к тому, что блокировка
будет удерживаться дольше, чем предполагалось, если мы невнимательны к времени
жизни MutexGuard<T>.
Код в листинге 21-20, который использует let job = receiver.lock().unwrap().recv().unwrap();, работает потому, что с let любые
временные значения, использованные в выражении справа от знака равенства,
немедленно сбрасываются, когда инструкция let заканчивается. Однако while let (а также if let и match) не сбрасывает временные значения до конца
соответствующего блока. В листинге 21-21 блокировка остается удержанной на все
время вызова job(), а значит, другие экземпляры Worker не могут получать
задания.