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

Передача данных между потоками с помощью передачи сообщений

Один все более популярный подход к обеспечению безопасной конкурентности – передача сообщений, при которой потоки или акторы общаются, отправляя друг другу сообщения с данными. Вот идея в виде лозунга из документации языка Go: “Не общайтесь, разделяя память; вместо этого разделяйте память, общаясь.”

Чтобы реализовать конкурентность с отправкой сообщений, стандартная библиотека Rust предоставляет реализацию каналов. Канал – это общее понятие в программировании, с помощью которого данные отправляются из одного потока в другой.

Канал в программировании можно представить как направленный водный канал, например ручей или реку. Если поместить в реку что-то вроде резиновой уточки, она поплывет вниз по течению к концу водного пути.

У канала есть две половины: передатчик и приемник. Половина-передатчик – это место выше по течению, где вы кладете резиновую уточку в реку, а половина-приемник – это место, куда резиновая уточка попадает ниже по течению. Одна часть вашего кода вызывает методы у передатчика с данными, которые нужно отправить, а другая часть проверяет принимающий конец на наличие прибывших сообщений. Канал считается закрытым, если либо половина передатчика, либо половина приемника была удалена.

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

Сначала в листинге 16-6 мы создадим канал, но ничего с ним делать не будем. Обратите внимание, что этот код пока не скомпилируется, потому что Rust не может определить, значения какого типа мы хотим отправлять по каналу.

Имя файла: src/main.rs
use std::sync::mpsc;

fn main() {
    let (tx, rx) = mpsc::channel();
}
Listing 16-6: Создание канала и присваивание двух его половин переменным tx и rx

Мы создаем новый канал с помощью функции mpsc::channel; mpsc означает multiple producer, single consumer – несколько производителей, один потребитель. Коротко говоря, способ реализации каналов в стандартной библиотеке Rust означает, что канал может иметь несколько отправляющих концов, которые производят значения, но только один принимающий конец, который потребляет эти значения. Представьте, что несколько ручьев сливаются в одну большую реку: все, что отправлено вниз по любому из ручьев, в конце окажется в одной реке. Пока мы начнем с одного производителя, но добавим нескольких производителей, когда этот пример заработает.

Функция mpsc::channel возвращает кортеж: первый элемент – отправляющий конец, передатчик, а второй элемент – принимающий конец, приемник. Сокращения tx и rx традиционно используются во многих областях для transmitter (передатчик) и receiver (приемник) соответственно, поэтому мы называем переменные именно так, чтобы обозначить каждый конец. Мы используем инструкцию let с шаблоном, который деструктурирует кортеж; использование шаблонов в инструкциях let и деструктуризацию мы обсудим в главе 19. Пока достаточно знать, что использование инструкции let таким способом – удобный подход для извлечения частей кортежа, возвращаемого mpsc::channel.

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

Имя файла: src/main.rs
use std::sync::mpsc;
use std::thread;

fn main() {
    let (tx, rx) = mpsc::channel();

    thread::spawn(move || {
        let val = String::from("hi");
        tx.send(val).unwrap();
    });
}
Listing 16-7: Перемещение tx в порожденный поток и отправка "hi"

Снова мы используем thread::spawn для создания нового потока, а затем используем move, чтобы переместить tx в замыкание, и порожденный поток стал владельцем tx. Порожденный поток должен владеть передатчиком, чтобы иметь возможность отправлять сообщения через канал.

У передатчика есть метод send, принимающий значение, которое мы хотим отправить. Метод send возвращает тип Result<T, E>, поэтому если приемник уже был удален и отправлять значение больше некуда, операция отправки вернет ошибку. В этом примере мы вызываем unwrap, чтобы вызвать панику в случае ошибки. Но в настоящем приложении мы обработали бы ее правильно: вернитесь к главе 9, чтобы вспомнить стратегии корректной обработки ошибок.

В листинге 16-8 мы получим значение из приемника в основном потоке. Это похоже на извлечение резиновой уточки из воды в конце реки или получение сообщения чата.

Имя файла: src/main.rs
use std::sync::mpsc;
use std::thread;

fn main() {
    let (tx, rx) = mpsc::channel();

    thread::spawn(move || {
        let val = String::from("hi");
        tx.send(val).unwrap();
    });

    let received = rx.recv().unwrap();
    println!("Got: {received}");
}
Listing 16-8: Получение значения "hi" в основном потоке и его печать

У приемника есть два полезных метода: recv и try_recv. Мы используем recv, сокращение от receive, который блокирует выполнение основного потока и ждет, пока значение не будет отправлено по каналу. Как только значение отправлено, recv вернет его в Result<T, E>. Когда передатчик закрывается, recv возвращает ошибку, сигнализируя, что больше значений не будет.

Метод try_recv не блокирует, а вместо этого сразу возвращает Result<T, E>: значение Ok, содержащее сообщение, если оно доступно, и значение Err, если сообщений в этот раз нет. Использование try_recv полезно, если этому потоку нужно выполнять другую работу, пока он ждет сообщений: мы могли бы написать цикл, который время от времени вызывает try_recv, обрабатывает сообщение, если оно доступно, а иначе некоторое время выполняет другую работу до следующей проверки.

В этом примере для простоты мы использовали recv; у основного потока нет другой работы, кроме ожидания сообщений, поэтому блокировка основного потока здесь уместна.

Когда мы запустим код из листинга 16-8, мы увидим значение, напечатанное из основного потока:

Got: hi

Отлично!

Передача владения через каналы

Правила владения играют жизненно важную роль в отправке сообщений, потому что они помогают писать безопасный конкурентный код. Предотвращение ошибок в конкурентном программировании – преимущество мышления о владении во всех ваших программах на Rust. Проведем эксперимент, чтобы показать, как каналы и владение работают вместе для предотвращения проблем: мы попробуем использовать значение val в порожденном потоке после того, как отправим его по каналу. Попробуйте скомпилировать код из листинга 16-9, чтобы увидеть, почему этот код не разрешен.

Имя файла: src/main.rs
use std::sync::mpsc;
use std::thread;

fn main() {
    let (tx, rx) = mpsc::channel();

    thread::spawn(move || {
        let val = String::from("hi");
        tx.send(val).unwrap();
        println!("val is {val}");
    });

    let received = rx.recv().unwrap();
    println!("Got: {received}");
}
Listing 16-9: Попытка использовать val после того, как мы отправили его по каналу

Здесь мы пытаемся напечатать val после того, как отправили его по каналу через tx.send. Разрешить это было бы плохой идеей: после отправки значения в другой поток этот поток мог бы изменить или удалить его до того, как мы попробуем использовать значение снова. Потенциально изменения, внесенные другим потоком, могли бы вызвать ошибки или неожиданные результаты из-за несогласованных или несуществующих данных. Однако Rust выдает ошибку, если мы пытаемся скомпилировать код из листинга 16-9:

$ cargo run
   Compiling message-passing v0.1.0 (file:///projects/message-passing)
error[E0382]: borrow of moved value: `val`
  --> src/main.rs:10:27
   |
 8 |         let val = String::from("hi");
   |             --- move occurs because `val` has type `String`, which does not implement the `Copy` trait
 9 |         tx.send(val).unwrap();
   |                 --- value moved here
10 |         println!("val is {val}");
   |                           ^^^ value borrowed here after move
   |
   = note: this error originates in the macro `$crate::format_args_nl` which comes from the expansion of the macro `println` (in Nightly builds, run with -Z macro-backtrace for more info)

For more information about this error, try `rustc --explain E0382`.
error: could not compile `message-passing` (bin "message-passing") due to 1 previous error

Наша ошибка конкурентности привела к ошибке времени компиляции. Функция send забирает владение своим параметром, и когда значение перемещается, приемник получает владение им. Это не дает нам случайно использовать значение снова после отправки; система владения проверяет, что все в порядке.

Отправка нескольких значений

Код в листинге 16-8 скомпилировался и выполнился, но он не показал нам ясно, что два отдельных потока разговаривают друг с другом через канал.

В листинге 16-10 мы внесли некоторые изменения, которые докажут, что код из листинга 16-8 выполняется конкурентно: теперь порожденный поток будет отправлять несколько сообщений и делать паузу в одну секунду между каждым сообщением.

Имя файла: src/main.rs
use std::sync::mpsc;
use std::thread;
use std::time::Duration;

fn main() {
    let (tx, rx) = mpsc::channel();

    thread::spawn(move || {
        let vals = vec![
            String::from("hi"),
            String::from("from"),
            String::from("the"),
            String::from("thread"),
        ];

        for val in vals {
            tx.send(val).unwrap();
            thread::sleep(Duration::from_secs(1));
        }
    });

    for received in rx {
        println!("Got: {received}");
    }
}
Listing 16-10: Отправка нескольких сообщений и пауза между каждым из них

На этот раз у порожденного потока есть вектор строк, которые мы хотим отправить в основной поток. Мы перебираем их, отправляя каждую отдельно, и делаем паузу между отправками, вызывая функцию thread::sleep со значением Duration в одну секунду.

В основном потоке мы больше не вызываем функцию recv явно: вместо этого мы обращаемся с rx как с итератором. Для каждого полученного значения мы печатаем его. Когда канал закрывается, итерация завершится.

При запуске кода из листинга 16-10 вы должны увидеть следующий вывод с секундной паузой между строками:

Got: hi
Got: from
Got: the
Got: thread

Поскольку у нас нет кода, который делает паузу или задержку в цикле for основного потока, мы можем сказать, что основной поток ждет получения значений из порожденного потока.

Создание нескольких производителей

Ранее мы упоминали, что mpsc – это сокращение от multiple producer, single consumer – несколько производителей, один потребитель. Используем mpsc и расширим код из листинга 16-10, чтобы создать несколько потоков, каждый из которых отправляет значения одному и тому же приемнику. Мы можем сделать это, клонировав передатчик, как показано в листинге 16-11.

Имя файла: src/main.rs
use std::sync::mpsc;
use std::thread;
use std::time::Duration;

fn main() {
    // --snip--

    let (tx, rx) = mpsc::channel();

    let tx1 = tx.clone();
    thread::spawn(move || {
        let vals = vec![
            String::from("hi"),
            String::from("from"),
            String::from("the"),
            String::from("thread"),
        ];

        for val in vals {
            tx1.send(val).unwrap();
            thread::sleep(Duration::from_secs(1));
        }
    });

    thread::spawn(move || {
        let vals = vec![
            String::from("more"),
            String::from("messages"),
            String::from("for"),
            String::from("you"),
        ];

        for val in vals {
            tx.send(val).unwrap();
            thread::sleep(Duration::from_secs(1));
        }
    });

    for received in rx {
        println!("Got: {received}");
    }

    // --snip--
}
Listing 16-11: Отправка нескольких сообщений от нескольких производителей

На этот раз перед созданием первого порожденного потока мы вызываем clone у передатчика. Это даст нам новый передатчик, который мы сможем передать первому порожденному потоку. Исходный передатчик мы передаем второму порожденному потоку. Так мы получаем два потока, каждый из которых отправляет разные сообщения одному приемнику.

Когда вы запустите код, ваш вывод должен выглядеть примерно так:

Got: hi
Got: more
Got: from
Got: messages
Got: for
Got: the
Got: thread
Got: you

Вы можете увидеть значения в другом порядке в зависимости от вашей системы. Именно это делает конкурентность одновременно интересной и трудной. Если поэкспериментировать с thread::sleep, задавая разные значения в разных потоках, каждый запуск станет более недетерминированным и каждый раз будет создавать разный вывод.

Теперь, когда мы посмотрели, как работают каналы, рассмотрим другой метод конкурентности.