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-10 реализует обработку запроса к /sleep с имитированным медленным ответом, из-за которого сервер заснет на пять секунд перед ответом.

Имя файла: src/main.rs
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();
}
Listing 21-10: Имитация медленного запроса с помощью сна на пять секунд

Теперь мы перешли с 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.

Имя файла: src/main.rs
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();
}
Listing 21-11: Создание нового потока для каждого потока данных

Как вы узнали в главе 16, thread::spawn создаст новый поток, а затем запустит код в замыкании в новом потоке. Если вы запустите этот код и загрузите /sleep в браузере, а затем / еще в двух вкладках браузера, вы действительно увидите, что запросам к / не нужно ждать завершения /sleep. Однако, как мы уже упоминали, со временем это перегрузит систему, потому что вы будете создавать новые потоки без какого-либо ограничения.

Возможно, вы также помните из главы 17, что это именно такая ситуация, где async и await действительно сияют! Держите это в уме, пока мы строим пул потоков, и подумайте, как все выглядело бы иначе или так же с async.

Создание конечного числа потоков

Мы хотим, чтобы наш пул потоков работал похожим и привычным образом, чтобы переход от потоков к пулу потоков не требовал больших изменений в коде, который использует наш API. Листинг 21-12 показывает гипотетический интерфейс для структуры ThreadPool, которую мы хотим использовать вместо thread::spawn.

Имя файла: src/main.rs
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();
}
Listing 21-12: Наш идеальный интерфейс 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, какое у нас может быть:

Имя файла: src/lib.rs
pub struct ThreadPool;

Затем отредактируйте файл main.rs, чтобы ввести ThreadPool из библиотечного крейта в область видимости, добавив следующий код в начало src/main.rs:

Имя файла: 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, которая будет обладать этими характеристиками:

Имя файла: src/lib.rs
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 с этими ограничениями:

Имя файла: src/lib.rs
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.

Имя файла: src/lib.rs
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,
    {
    }
}
Listing 21-13: Реализация 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, содержащий их.

Имя файла: src/lib.rs
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,
    {
    }
}
Listing 21-14: Создание вектора для хранения потоков в 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 таким образом:

  1. Определить структуру Worker, которая хранит id и JoinHandle<()>.
  2. Изменить ThreadPool, чтобы он хранил вектор экземпляров Worker.
  3. Определить функцию Worker::new, которая принимает число id и возвращает экземпляр Worker, хранящий id и поток, созданный с пустым замыканием.
  4. В ThreadPool::new использовать счетчик цикла for, чтобы сгенерировать id, создать новый Worker с этим id и сохранить Worker в векторе.

Если вы готовы к испытанию, попробуйте реализовать эти изменения самостоятельно перед тем, как смотреть код в листинге 21-15.

Готовы? Вот листинг 21-15 с одним из способов внести предыдущие изменения.

Имя файла: src/lib.rs
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 }
    }
}
Listing 21-15: Изменение 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, которые отправят задание своему потоку. Вот план:

  1. ThreadPool создаст канал и сохранит отправителя.
  2. Каждый Worker сохранит получателя.
  3. Мы создадим новую структуру Job, которая будет хранить замыкания, которые мы хотим отправлять по каналу.
  4. Метод execute отправит задание, которое хочет выполнить, через отправителя.
  5. В своем потоке Worker будет циклически получать данные от своего получателя и выполнять замыкания любых полученных заданий.

Начнем с создания канала в ThreadPool::new и хранения отправителя в экземпляре ThreadPool, как показано в листинге 21-16. Структура Job пока ничего не хранит, но будет типом элементов, которые мы отправляем по каналу.

Имя файла: src/lib.rs
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 }
    }
}
Listing 21-16: Хранение в ThreadPool отправителя канала для экземпляров Job

В ThreadPool::new мы создаем новый канал и заставляем пул хранить отправителя. Это успешно скомпилируется.

Попробуем передать получателя канала каждому Worker во время создания канала пулом потоков. Мы знаем, что хотим использовать получателя в потоке, который запускают экземпляры Worker, поэтому сошлемся на параметр receiver в замыкании. Код в листинге 21-17 пока не совсем скомпилируется.

Имя файла: src/lib.rs
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 }
    }
}
Listing 21-17: Передача получателя каждому 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 показывает изменения, которые нужно внести.

Имя файла: src/lib.rs
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 }
    }
}
Listing 21-18: Совместное использование получателя через Arc и Mutex

В ThreadPool::new мы помещаем получателя в Arc и Mutex. Для каждого нового Worker мы клонируем Arc, чтобы увеличить счетчик ссылок, благодаря чему экземпляры Worker могут совместно владеть получателем.

С этими изменениями код компилируется! Мы приближаемся!

Реализация метода execute

Наконец реализуем метод execute у ThreadPool. Мы также изменим Job со структуры на псевдоним типа для трейт-объекта, который хранит тип замыкания, получаемого execute. Как обсуждалось в разделе «Синонимы типов и псевдонимы типов» главы 20, псевдонимы типов позволяют нам сокращать длинные типы для удобства использования. Посмотрите на листинг 21-19.

Имя файла: src/lib.rs
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 }
    }
}
Listing 21-19: Создание псевдонима Job для Box с замыканием и отправка задания по каналу

После создания нового экземпляра Job с помощью замыкания, полученного в execute, мы отправляем это задание через отправляющий конец канала. Мы вызываем unwrap у send на случай, если отправка завершится неудачно. Это может произойти, например, если мы остановим выполнение всех наших потоков, то есть принимающий конец перестанет получать новые сообщения. Сейчас мы не можем остановить выполнение потоков: наши потоки продолжают выполняться, пока существует пул. Причина, по которой мы используем unwrap, в том, что мы знаем: случай ошибки не произойдет, но компилятор этого не знает.

Но мы еще не совсем закончили! В Worker замыкание, передаваемое thread::spawn, все еще только ссылается на принимающий конец канала. Вместо этого нам нужно, чтобы замыкание выполняло бесконечный цикл, запрашивая у принимающего конца канала задание и выполняя его, когда оно получено. Внесем изменение в Worker::new, показанное в листинге 21-20.

Имя файла: 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();
    }
}

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 }
    }
}
Listing 21-20: Получение и выполнение заданий в потоке экземпляра 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.

Имя файла: 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();
    }
}

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 }
    }
}
Listing 21-21: Альтернативная реализация 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 не могут получать задания.