Применение конкурентности с 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.
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;
}
});
}
В качестве отправной точки мы настраиваем функцию 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, поэтому
после ожидания мы также разворачиваем его.
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();
});
}
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. Вывод мы игнорируем, потому что это всего лишь кортеж,
содержащий два единичных значения.
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;
});
}
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-блока — не порождая отдельную задачу, как мы порождали отдельный поток.
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}'");
});
}
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.
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}'");
}
});
}
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 последовательно, мы просто вернулись бы к
последовательному потоку выполнения — именно к тому, чего мы пытаемся не
делать.
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;
});
}
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.
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;
});
}
Когда мы запускаем эту версию кода, она корректно завершает работу после отправки и получения последнего сообщения. Далее посмотрим, что нужно изменить, чтобы отправлять данные из более чем одного future.
Объединение нескольких futures с помощью макроса join!
Этот async-канал также является каналом с несколькими производителями, поэтому
мы можем вызвать clone у tx, если хотим отправлять сообщения из нескольких
futures, как показано в листинге 17-13.
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);
});
}
Сначала мы клонируем 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. Далее обсудим, как и почему сообщать среде выполнения, что она может переключиться на другую задачу.