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

Применение конкурентности с async

В этом разделе мы применим async к некоторым из тех же задач конкурентности, которые решали с потоками в главе 16. Поскольку там мы уже обсудили многие ключевые идеи, в этом разделе сосредоточимся на том, чем потоки отличаются от futures.

Во многих случаях API для работы с конкурентностью с использованием async очень похожи на API для использования потоков. В других случаях они в итоге оказываются совсем другими. Даже когда API для потоков и async выглядят похожими, они часто ведут себя по-разному — и почти всегда имеют разные характеристики производительности.

Создание новой задачи с помощью spawn_task

Первой операцией, которую мы разбирали в разделе «Создание нового потока с помощью spawn» главы 16, был счет в двух отдельных потоках. Давайте сделаем то же самое с помощью async. Крейт trpl предоставляет функцию spawn_task, которая выглядит очень похожей на API thread::spawn, и функцию sleep, которая является async-версией API thread::sleep. Мы можем использовать их вместе, чтобы реализовать пример со счетом, как показано в листинге 17-6.

Имя файла: src/main.rs
extern crate trpl; // required for mdbook test

use std::time::Duration;

fn main() {
    trpl::block_on(async {
        trpl::spawn_task(async {
            for i in 1..10 {
                println!("hi number {i} from the first task!");
                trpl::sleep(Duration::from_millis(500)).await;
            }
        });

        for i in 1..5 {
            println!("hi number {i} from the second task!");
            trpl::sleep(Duration::from_millis(500)).await;
        }
    });
}
Listing 17-6: Создание новой задачи, чтобы печатать одно, пока основная задача печатает другое

В качестве отправной точки мы настраиваем функцию main с помощью trpl::block_on, чтобы наша функция верхнего уровня могла быть async.

Примечание: Начиная с этого места в главе, каждый пример будет включать точно такой же оборачивающий код с trpl::block_on в main, поэтому мы часто будем пропускать его так же, как пропускаем main. Не забывайте включать его в свой код!

Затем внутри этого блока мы пишем два цикла, каждый из которых содержит вызов trpl::sleep, ожидающий полсекунды (500 миллисекунд) перед отправкой следующего сообщения. Один цикл мы помещаем в тело trpl::spawn_task, а другой — в цикл for верхнего уровня. Мы также добавляем await после вызовов sleep.

Этот код ведет себя похоже на реализацию на основе потоков, включая то, что при запуске вы можете увидеть сообщения в своем терминале в другом порядке:

hi number 1 from the second task!
hi number 1 from the first task!
hi number 2 from the first task!
hi number 2 from the second task!
hi number 3 from the first task!
hi number 3 from the second task!
hi number 4 from the first task!
hi number 4 from the second task!
hi number 5 from the first task!

Эта версия останавливается, как только цикл for в теле основного async-блока завершается, потому что задача, порожденная spawn_task, завершается принудительно, когда завершается функция main. Если вы хотите, чтобы она выполнилась до конца задачи, нужно использовать дескриптор соединения, чтобы дождаться завершения первой задачи. С потоками мы использовали метод join, чтобы «блокироваться» до завершения выполнения потока. В листинге 17-7 мы можем использовать await, чтобы сделать то же самое, потому что сам дескриптор задачи является future. Его тип Output — это Result, поэтому после ожидания мы также разворачиваем его.

Имя файла: src/main.rs
extern crate trpl; // required for mdbook test

use std::time::Duration;

fn main() {
    trpl::block_on(async {
        let handle = trpl::spawn_task(async {
            for i in 1..10 {
                println!("hi number {i} from the first task!");
                trpl::sleep(Duration::from_millis(500)).await;
            }
        });

        for i in 1..5 {
            println!("hi number {i} from the second task!");
            trpl::sleep(Duration::from_millis(500)).await;
        }

        handle.await.unwrap();
    });
}
Listing 17-7: Использование await с дескриптором соединения, чтобы выполнить задачу до конца

Эта обновленная версия выполняется до тех пор, пока не завершатся оба цикла:

hi number 1 from the second task!
hi number 1 from the first task!
hi number 2 from the first task!
hi number 2 from the second task!
hi number 3 from the first task!
hi number 3 from the second task!
hi number 4 from the first task!
hi number 4 from the second task!
hi number 5 from the first task!
hi number 6 from the first task!
hi number 7 from the first task!
hi number 8 from the first task!
hi number 9 from the first task!

Пока что кажется, что async и потоки дают похожие результаты, только с разным синтаксисом: мы используем await вместо вызова join у дескриптора соединения и ожидаем вызовы sleep.

Более существенное отличие в том, что для этого нам не пришлось порождать еще один поток операционной системы. На самом деле здесь даже не нужно порождать задачу. Поскольку async-блоки компилируются в анонимные futures, мы можем поместить каждый цикл в async-блок и позволить среде выполнения выполнить оба до завершения с помощью функции trpl::join.

В разделе «Ожидание завершения всех потоков» главы 16 мы показали, как использовать метод join у типа JoinHandle, который возвращается при вызове std::thread::spawn. Функция trpl::join похожа на него, но предназначена для futures. Когда вы передаете ей два future, она создает один новый future, вывод которого является кортежем, содержащим вывод каждого переданного future после того, как оба завершатся. Таким образом, в листинге 17-8 мы используем trpl::join, чтобы дождаться завершения и fut1, и fut2. Мы ожидаем не fut1 и fut2, а новый future, созданный trpl::join. Вывод мы игнорируем, потому что это всего лишь кортеж, содержащий два единичных значения.

Имя файла: src/main.rs
extern crate trpl; // required for mdbook test

use std::time::Duration;

fn main() {
    trpl::block_on(async {
        let fut1 = async {
            for i in 1..10 {
                println!("hi number {i} from the first task!");
                trpl::sleep(Duration::from_millis(500)).await;
            }
        };

        let fut2 = async {
            for i in 1..5 {
                println!("hi number {i} from the second task!");
                trpl::sleep(Duration::from_millis(500)).await;
            }
        };

        trpl::join(fut1, fut2).await;
    });
}
Listing 17-8: Использование trpl::join для ожидания двух анонимных futures

Когда мы запускаем это, мы видим, что оба future выполняются до завершения:

hi number 1 from the first task!
hi number 1 from the second task!
hi number 2 from the first task!
hi number 2 from the second task!
hi number 3 from the first task!
hi number 3 from the second task!
hi number 4 from the first task!
hi number 4 from the second task!
hi number 5 from the first task!
hi number 6 from the first task!
hi number 7 from the first task!
hi number 8 from the first task!
hi number 9 from the first task!

Теперь вы будете видеть один и тот же порядок каждый раз, что сильно отличается от того, что мы видели с потоками и с trpl::spawn_task в листинге 17-7. Причина в том, что функция trpl::join является справедливой: она проверяет каждый future одинаково часто, чередуясь между ними, и никогда не позволяет одному сильно вырваться вперед, если другой тоже готов. В случае с потоками операционная система решает, какой поток проверять и как долго позволять ему выполняться. В async Rust среда выполнения решает, какую задачу проверять. (На практике детали становятся сложнее, потому что async-среда выполнения может использовать потоки операционной системы под капотом как часть того, как она управляет конкурентностью, поэтому гарантия справедливости может требовать от среды выполнения больше работы, но это все равно возможно!) Среды выполнения не обязаны гарантировать справедливость для каждой конкретной операции, и они часто предлагают разные API, чтобы вы могли выбрать, нужна ли вам справедливость или нет.

Попробуйте некоторые из этих вариантов ожидания futures и посмотрите, что они делают:

  • Уберите async-блок вокруг одного или обоих циклов.
  • Ожидайте каждый async-блок сразу после его определения.
  • Оберните только первый цикл в async-блок и ожидайте получившийся future после тела второго цикла.

В качестве дополнительной задачи попробуйте понять, каким будет вывод в каждом случае, до запуска кода!

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

Совместное использование данных между futures тоже будет знакомым: мы снова используем передачу сообщений, но на этот раз с async-версиями типов и функций. Мы пойдем немного другим путем, чем в разделе «Передача данных между потоками с помощью передачи сообщений» главы 16, чтобы показать некоторые ключевые различия между конкурентностью на основе потоков и конкурентностью на основе futures. В листинге 17-9 мы начнем всего с одного async-блока — не порождая отдельную задачу, как мы порождали отдельный поток.

Имя файла: src/main.rs
extern crate trpl; // required for mdbook test

fn main() {
    trpl::block_on(async {
        let (tx, mut rx) = trpl::channel();

        let val = String::from("hi");
        tx.send(val).unwrap();

        let received = rx.recv().await.unwrap();
        println!("received '{received}'");
    });
}
Listing 17-9: Создание async-канала и присваивание двух половин tx и rx

Здесь мы используем trpl::channel, async-версию API канала с несколькими производителями и одним потребителем, который мы использовали с потоками еще в главе 16. Async-версия API лишь немного отличается от версии на основе потоков: она использует изменяемый, а не неизменяемый приемник rx, и ее метод recv создает future, который нужно ожидать, вместо того чтобы возвращать значение напрямую. Теперь мы можем отправлять сообщения от отправителя к приемнику. Обратите внимание, что нам не нужно порождать отдельный поток или даже задачу; нам нужно лишь ожидать вызов rx.recv.

Синхронный метод Receiver::recv в std::mpsc::channel блокируется, пока не получит сообщение. Метод trpl::Receiver::recv этого не делает, потому что он async. Вместо блокировки он возвращает управление среде выполнения, пока либо не будет получено сообщение, либо не закроется отправляющая сторона канала. В противоположность этому мы не ожидаем вызов send, потому что он не блокируется. Ему и не нужно блокироваться, потому что канал, в который мы отправляем, неограниченный.

Примечание: Поскольку весь этот async-код выполняется в async-блоке внутри вызова trpl::block_on, все внутри него может избегать блокировки. Однако код снаружи будет блокироваться на возврате функции block_on. В этом и состоит весь смысл функции trpl::block_on: она позволяет вам выбрать, где блокироваться на некотором наборе async-кода, а значит, где переходить между синхронным и асинхронным кодом.

Обратите внимание на две вещи в этом примере. Во-первых, сообщение придет сразу. Во-вторых, хотя здесь мы используем future, конкурентности пока нет. Все в листинге происходит последовательно, точно так же, как происходило бы, если бы futures вообще не были задействованы.

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

Имя файла: src/main.rs
extern crate trpl; // required for mdbook test

use std::time::Duration;

fn main() {
    trpl::block_on(async {
        let (tx, mut rx) = trpl::channel();

        let vals = vec![
            String::from("hi"),
            String::from("from"),
            String::from("the"),
            String::from("future"),
        ];

        for val in vals {
            tx.send(val).unwrap();
            trpl::sleep(Duration::from_millis(500)).await;
        }

        while let Some(value) = rx.recv().await {
            println!("received '{value}'");
        }
    });
}
Listing 17-10: Отправка и получение нескольких сообщений через async-канал и сон с await между каждым сообщением

Помимо отправки сообщений, нам нужно их получать. В этом случае, поскольку мы знаем, сколько сообщений поступит, мы могли бы сделать это вручную, вызвав rx.recv().await четыре раза. Однако в реальном мире мы обычно будем ждать некоторое неизвестное количество сообщений, поэтому нужно продолжать ждать, пока мы не определим, что сообщений больше нет.

В листинге 16-10 мы использовали цикл for, чтобы обработать все элементы, полученные из синхронного канала. Однако в Rust пока нет способа использовать цикл for с асинхронно производимой серией элементов, поэтому нам нужно использовать цикл, которого мы раньше не видели: условный цикл while let. Это циклическая версия конструкции if let, которую мы видели в разделе «Краткий поток управления с if let и let...else» главы 6. Цикл продолжит выполняться, пока заданный им шаблон продолжает соответствовать значению.

Вызов rx.recv создает future, который мы ожидаем. Среда выполнения приостановит future, пока он не будет готов. Когда приходит сообщение, future разрешается в Some(message) столько раз, сколько приходят сообщения. Когда канал закрывается, независимо от того, пришли ли какие-либо сообщения, future вместо этого разрешается в None, чтобы указать, что значений больше нет и поэтому мы должны прекратить опрос, то есть прекратить ожидание.

Цикл while let собирает все это вместе. Если результат вызова rx.recv().await — это Some(message), мы получаем доступ к сообщению и можем использовать его в теле цикла, как могли бы с if let. Если результат — None, цикл заканчивается. Каждый раз, когда цикл завершается, он снова доходит до точки await, поэтому среда выполнения снова приостанавливает его до появления другого сообщения.

Теперь код успешно отправляет и получает все сообщения. К сожалению, все еще есть пара проблем. Во-первых, сообщения не приходят с интервалами в полсекунды. Они приходят все сразу, через 2 секунды (2 000 миллисекунд) после запуска программы. Во-вторых, эта программа также никогда не завершается! Вместо этого она вечно ждет новые сообщения. Вам нужно будет остановить ее с помощью ctrl-C.

Код внутри одного async-блока выполняется линейно

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

В листинге 17-10 есть только один async-блок, поэтому все внутри него выполняется линейно. Конкурентности по-прежнему нет. Все вызовы tx.send происходят вперемешку со всеми вызовами trpl::sleep и связанными с ними точками await. Только после этого цикл while let получает возможность пройти через какие-либо точки await на вызовах recv.

Чтобы получить нужное поведение, при котором задержка сна происходит между каждым сообщением, нужно поместить операции tx и rx в их собственные async-блоки, как показано в листинге 17-11. Тогда среда выполнения сможет исполнять каждый из них отдельно с помощью trpl::join, как в листинге 17-8. И снова мы ожидаем результат вызова trpl::join, а не отдельные futures. Если бы мы ожидали отдельные futures последовательно, мы просто вернулись бы к последовательному потоку выполнения — именно к тому, чего мы пытаемся не делать.

Имя файла: src/main.rs
extern crate trpl; // required for mdbook test

use std::time::Duration;

fn main() {
    trpl::block_on(async {
        let (tx, mut rx) = trpl::channel();

        let tx_fut = async {
            let vals = vec![
                String::from("hi"),
                String::from("from"),
                String::from("the"),
                String::from("future"),
            ];

            for val in vals {
                tx.send(val).unwrap();
                trpl::sleep(Duration::from_millis(500)).await;
            }
        };

        let rx_fut = async {
            while let Some(value) = rx.recv().await {
                println!("received '{value}'");
            }
        };

        trpl::join(tx_fut, rx_fut).await;
    });
}
Listing 17-11: Разделение send и recv по собственным блокам async и ожидание futures этих блоков

С обновленным кодом в листинге 17-11 сообщения печатаются с интервалами в 500 миллисекунд, а не все сразу после 2 секунд.

Перемещение владения в async-блок

Однако программа все еще никогда не завершается из-за того, как цикл while let взаимодействует с trpl::join:

  • Future, возвращенный из trpl::join, завершается только после того, как завершатся оба future, переданных в него.
  • Future tx_fut завершается после того, как заканчивается его сон после отправки последнего сообщения в vals.
  • Future rx_fut не завершится, пока не закончится цикл while let.
  • Цикл while let не закончится, пока ожидание rx.recv не даст None.
  • Ожидание rx.recv вернет None только после того, как другой конец канала будет закрыт.
  • Канал закроется только если мы вызовем rx.close или когда отправляющая сторона, tx, будет отброшена.
  • Мы нигде не вызываем rx.close, а tx не будет отброшен до тех пор, пока не завершится самый внешний async-блок, переданный в trpl::block_on.
  • Блок не может завершиться, потому что он заблокирован на завершении trpl::join, что возвращает нас к началу этого списка.

Сейчас async-блок, в котором мы отправляем сообщения, только заимствует tx, потому что отправка сообщения не требует владения. Но если бы мы могли переместить tx в этот async-блок, он был бы отброшен после завершения блока. В разделе «Захват ссылок или перемещение владения» главы 13 вы узнали, как использовать ключевое слово move с замыканиями, и, как обсуждалось в разделе «Использование замыканий move с потоками» главы 16, нам часто нужно перемещать данные в замыкания при работе с потоками. Та же базовая динамика применима к async-блокам, поэтому ключевое слово move работает с async-блоками так же, как с замыканиями.

В листинге 17-12 мы меняем блок, используемый для отправки сообщений, с async на async move.

Имя файла: src/main.rs
extern crate trpl; // required for mdbook test

use std::time::Duration;

fn main() {
    trpl::block_on(async {
        let (tx, mut rx) = trpl::channel();

        let tx_fut = async move {
            // --snip--
            let vals = vec![
                String::from("hi"),
                String::from("from"),
                String::from("the"),
                String::from("future"),
            ];

            for val in vals {
                tx.send(val).unwrap();
                trpl::sleep(Duration::from_millis(500)).await;
            }
        };

        let rx_fut = async {
            while let Some(value) = rx.recv().await {
                println!("received '{value}'");
            }
        };

        trpl::join(tx_fut, rx_fut).await;
    });
}
Listing 17-12: Переработка кода из листинга 17-11, которая корректно завершает работу после выполнения

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

Объединение нескольких futures с помощью макроса join!

Этот async-канал также является каналом с несколькими производителями, поэтому мы можем вызвать clone у tx, если хотим отправлять сообщения из нескольких futures, как показано в листинге 17-13.

Имя файла: src/main.rs
extern crate trpl; // required for mdbook test

use std::time::Duration;

fn main() {
    trpl::block_on(async {
        let (tx, mut rx) = trpl::channel();

        let tx1 = tx.clone();
        let tx1_fut = async move {
            let vals = vec![
                String::from("hi"),
                String::from("from"),
                String::from("the"),
                String::from("future"),
            ];

            for val in vals {
                tx1.send(val).unwrap();
                trpl::sleep(Duration::from_millis(500)).await;
            }
        };

        let rx_fut = async {
            while let Some(value) = rx.recv().await {
                println!("received '{value}'");
            }
        };

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

            for val in vals {
                tx.send(val).unwrap();
                trpl::sleep(Duration::from_millis(1500)).await;
            }
        };

        trpl::join!(tx1_fut, tx_fut, rx_fut);
    });
}
Listing 17-13: Использование нескольких производителей с async-блоками

Сначала мы клонируем tx, создавая tx1 вне первого async-блока. Мы перемещаем tx1 в этот блок точно так же, как раньше поступили с tx. Затем, позже, мы перемещаем исходный tx в новый async-блок, где отправляем больше сообщений с немного более медленной задержкой. Так получилось, что мы помещаем этот новый async-блок после async-блока для получения сообщений, но он с тем же успехом мог бы идти перед ним. Главное — порядок, в котором futures ожидаются, а не порядок, в котором они создаются.

Оба async-блока для отправки сообщений должны быть блоками async move, чтобы и tx, и tx1 были отброшены при завершении этих блоков. Иначе мы снова окажемся в том же бесконечном цикле, с которого начали.

Наконец, мы переключаемся с trpl::join на trpl::join!, чтобы обработать дополнительный future: макрос join! ожидает произвольное число futures, когда количество futures известно во время компиляции. Ожидание коллекции с неизвестным числом futures мы обсудим позже в этой главе.

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

received 'hi'
received 'more'
received 'from'
received 'the'
received 'messages'
received 'future'
received 'for'
received 'you'

Мы исследовали, как использовать передачу сообщений для отправки данных между futures, как код внутри async-блока выполняется последовательно, как перемещать владение в async-блок и как объединять несколько futures. Далее обсудим, как и почему сообщать среде выполнения, что она может переключиться на другую задачу.