Макрос select
macro_rules! select {
{
$(
biased;
)?
$(
$bind:pat = $fut:expr $(, if $cond:expr)? => $handler:expr,
)*
$(
else => $els:expr $(,)?
)?
} => { ... };
}
macros.Ожидает завершения нескольких параллельных ветвей, возвращая управление, когда завершается первая ветвь, и отменяя остальные.
Макрос select! должен использоваться внутри асинхронных функций, замыканий и блоков.
Макрос select! принимает одну или несколько ветвей следующего вида:
<pattern> = <async expression> (, if <precondition>)? => <handler>,Кроме того, макрос select! может содержать одну необязательную ветвь else, которая выполняется, если ни одна из других ветвей не соответствует своим шаблонам:
else => <expression>Макрос объединяет все выражения <async expression> и выполняет их параллельно в рамках текущей задачи. Как только первое выражение завершается значением, соответствующим его <pattern>, макрос select! возвращает результат вычисления выражения <handler> завершившейся ветви.
Кроме того, каждая ветвь может содержать необязательное предварительное условие if. Если предварительное условие возвращает false, ветвь отключается. Указанное <async expression> всё равно вычисляется, но полученное будущее никогда не опрашивается. Это удобно при использовании select! в цикле.
Полный жизненный цикл выражения select! выглядит следующим образом:
- Вычислить все указанные выражения
<precondition>. Если предварительное условие возвращаетfalse, отключить ветвь до конца текущего вызоваselect!. Повторный вход вselect!из-за цикла сбрасывает состояние «отключена». - Объединить
<async expression>всех ветвей, включая отключённые. Если ветвь отключена,<async expression>всё равно вычисляется, но полученное будущее не опрашивается. - Если отключены все ветви, перейти к шагу 6.
- Ожидать параллельно результаты всех оставшихся
<async expression>. - Когда
<async expression>возвращает значение, попытаться сопоставить его с указанным<pattern>. Если шаблон совпадает, вычислить<handler>и выполнить возврат. Если шаблон не совпадает, отключить текущую ветвь до конца текущего вызоваselect!. Продолжить с шага 3. - Вычислить выражение
else. Если выражение else не указано, вызвать панику.
Особенности выполнения
Поскольку все асинхронные выражения выполняются в текущей задаче, они могут выполняться параллельно, но не одновременно. Это означает, что все выражения выполняются в одном потоке, и если одна ветвь блокирует поток, остальные выражения не смогут продолжить выполнение. Если требуется параллельное выполнение, запустите каждое асинхронное выражение с помощью tokio::spawn и передайте дескриптор присоединения в select!.
Справедливость
По умолчанию select! случайным образом выбирает ветвь, которую проверит первой. Это обеспечивает некоторую степень справедливости при вызове select! в цикле с ветвями, которые всегда готовы к выполнению.
Это поведение можно изменить, добавив biased; в начало использования макроса. Подробности см. в примерах. В этом случае select будет опрашивать будущие в порядке их появления сверху вниз. Это может быть нужно по нескольким причинам:
- Генерация случайных чисел в
tokio::select!требует ненулевых затрат процессорного времени. - Будущие могут взаимодействовать таким образом, что известный порядок опроса имеет значение.
Однако у этого режима есть важный нюанс. Вы сами должны следить за справедливостью порядка опроса будущих. Например, если вы выбираете между потоком и будущим завершения работы, а поток содержит огромное количество сообщений с нулевым или почти нулевым интервалом между ними, следует поместить будущее завершения работы раньше в списке select!, чтобы оно всегда опрашивалось и не игнорировалось из-за того, что поток постоянно готов.
Паники
Макрос select! вызывает панику, если отключены все ветви и ветвь else не указана. Ветвь отключается, если указанное предварительное условие if возвращает false или если шаблон не соответствует результату <async expression>.
Безопасность отмены
При использовании select! в цикле для получения сообщений из нескольких источников убедитесь, что вызов получения безопасен при отмене, чтобы избежать потери сообщений. В этом разделе рассматриваются различные распространённые методы и описывается, безопасны ли они при отмене. Списки в этом разделе не являются исчерпывающими.
Безопасность отмены описывает, что происходит, когда будущее отбрасывается до завершения. Безопасность отмены зависит от поведения будущего, переданного в select!. Такое будущее может быть получено из асинхронного метода, асинхронного выражения или другой операции, создающей будущее.
Следующие методы безопасны при отмене:
tokio::sync::mpsc::Receiver::recvtokio::sync::mpsc::UnboundedReceiver::recvtokio::sync::broadcast::Receiver::recvtokio::sync::watch::Receiver::changedtokio::net::TcpListener::accepttokio::net::UnixListener::accepttokio::signal::unix::Signal::recv-
tokio::io::AsyncReadExt::readдля любогоAsyncRead -
tokio::io::AsyncReadExt::read_bufдля любогоAsyncRead -
tokio::io::AsyncWriteExt::writeдля любогоAsyncWrite -
tokio::io::AsyncWriteExt::write_bufдля любогоAsyncWrite -
tokio_stream::StreamExt::nextдля любогоStream -
futures::stream::StreamExt::nextдля любогоStream
Следующие методы небезопасны при отмене и могут привести к потере данных:
tokio::io::AsyncReadExt::read_exacttokio::io::AsyncReadExt::read_to_endtokio::io::AsyncReadExt::read_to_stringtokio::io::AsyncWriteExt::write_all
Следующие методы небезопасны при отмене, поскольку для обеспечения справедливости они используют очередь, а при отмене вы теряете своё место в очереди:
tokio::sync::Mutex::locktokio::sync::RwLock::readtokio::sync::RwLock::writetokio::sync::Semaphore::acquiretokio::sync::Notify::notified
Чтобы определить, безопасны ли при отмене ваши собственные методы, найдите места использования .await. Это связано с тем, что отмена асинхронного метода всегда происходит в точке .await. Если функция работает корректно, даже когда её выполнение начинается заново во время ожидания в точке .await, значит, она безопасна при отмене.
Безопасность отмены можно определить следующим образом: если у вас есть ещё не завершившееся будущее, его отбрасывание и создание заново не должны иметь никаких побочных эффектов. Это определение обусловлено ситуацией, когда select! используется в цикле. Без этой гарантии вы потеряете достигнутый прогресс, когда завершится другая ветвь и вы заново запустите select!, перейдя на следующую итерацию цикла.
Помните, что отмена операции, небезопасной при отмене, не обязательно является ошибкой. Например, если вы отменяете задачу из-за завершения работы приложения, вероятно, вас не беспокоит потеря частично прочитанных данных.
Примеры
Простой select с двумя ветвями.
async fn do_stuff_async() {
// async work
}
async fn more_async_work() {
// more here
}
tokio::select! {
_ = do_stuff_async() => {
println!("do_stuff_async() completed first")
}
_ = more_async_work() => {
println!("more_async_work() completed first")
}
};Простой выбор из потоков.
use tokio_stream::{self as stream, StreamExt};
let mut stream1 = stream::iter(vec![1, 2, 3]);
let mut stream2 = stream::iter(vec![4, 5, 6]);
let next = tokio::select! {
v = stream1.next() => v.unwrap(),
v = stream2.next() => v.unwrap(),
};
assert!(next == 1 || next == 4);Собрать содержимое двух потоков. В этом примере мы полагаемся на сопоставление с шаблоном и на то, что stream::iter является «слитым», то есть после завершения потока все вызовы next() возвращают None.
use tokio_stream::{self as stream, StreamExt};
let mut stream1 = stream::iter(vec![1, 2, 3]);
let mut stream2 = stream::iter(vec![4, 5, 6]);
let mut values = vec![];
loop {
tokio::select! {
Some(v) = stream1.next() => values.push(v),
Some(v) = stream2.next() => values.push(v),
else => break,
}
}
values.sort();
assert_eq!(&[1, 2, 3, 4, 5, 6], &values[..]);Одно и то же будущее можно использовать в нескольких выражениях select!, передав ссылку на него. Для этого будущее должно реализовывать Unpin. Будущее можно сделать Unpin, используя Box::pin или закрепив его в стеке.
В этом примере поток обрабатывается не более 1 секунды.
use tokio_stream::{self as stream, StreamExt};
use tokio::time::{self, Duration};
let mut stream = stream::iter(vec![1, 2, 3]);
let sleep = time::sleep(Duration::from_secs(1));
tokio::pin!(sleep);
loop {
tokio::select! {
maybe_v = stream.next() => {
if let Some(v) = maybe_v {
println!("got = {}", v);
} else {
break;
}
}
_ = &mut sleep => {
println!("timeout");
break;
}
}
}Объединение двух значений с помощью select!.
use tokio::sync::oneshot;
let (tx1, mut rx1) = oneshot::channel();
let (tx2, mut rx2) = oneshot::channel();
tokio::spawn(async move {
tx1.send("first").unwrap();
});
tokio::spawn(async move {
tx2.send("second").unwrap();
});
let mut a = None;
let mut b = None;
while a.is_none() || b.is_none() {
tokio::select! {
v1 = (&mut rx1), if a.is_none() => a = Some(v1.unwrap()),
v2 = (&mut rx2), if b.is_none() => b = Some(v2.unwrap()),
}
}
let res = (a.unwrap(), b.unwrap());
assert_eq!(res.0, "first");
assert_eq!(res.1, "second");Использование режима biased; для управления порядком опроса.
let mut count = 0u8;
loop {
tokio::select! {
// If you run this example without `biased;`, the polling order is
// pseudo-random, and the assertions on the value of count will
// (probably) fail.
biased;
_ = async {}, if count < 1 => {
count += 1;
assert_eq!(count, 1);
}
_ = async {}, if count < 2 => {
count += 1;
assert_eq!(count, 2);
}
_ = async {}, if count < 3 => {
count += 1;
assert_eq!(count, 3);
}
_ = async {}, if count < 4 => {
count += 1;
assert_eq!(count, 4);
}
else => {
break;
}
};
}Избегайте гонок при использовании предварительных условий if
Поскольку предварительные условия if используются для отключения ветвей select!, необходимо соблюдать осторожность, чтобы не пропустить значения.
Например, ниже показано неправильное использование sleep с if. Цель — многократно выполнять асинхронную задачу в течение не более 50 миллисекунд. Однако существует риск пропустить завершение sleep.
use tokio::time::{self, Duration};
async fn some_async_work() {
// do work
}
let sleep = time::sleep(Duration::from_millis(50));
tokio::pin!(sleep);
while !sleep.is_elapsed() {
tokio::select! {
_ = &mut sleep, if !sleep.is_elapsed() => {
println!("operation timed out");
}
_ = some_async_work() => {
println!("operation completed");
}
}
}
panic!("This example shows how not to do it!");
В приведённом выше примере sleep.is_elapsed() может вернуть true, даже если sleep.poll() никогда не возвращал Ready. Это создаёт потенциальное состояние гонки: срок действия sleep может истечь между проверкой while !sleep.is_elapsed() и вызовом select!, в результате чего вызов some_async_work() выполнится без прерывания, хотя время ожидания уже истекло.
Один из способов переписать приведённый выше пример без гонки:
use tokio::time::{self, Duration};
async fn some_async_work() {
// do work
}
let sleep = time::sleep(Duration::from_millis(50));
tokio::pin!(sleep);
loop {
tokio::select! {
_ = &mut sleep => {
println!("operation timed out");
break;
}
_ = some_async_work() => {
println!("operation completed");
}
}
}Альтернативы в экосистеме
Макрос select! — мощный инструмент для управления несколькими асинхронными ветвями, позволяющий выполнять задачи параллельно в одном потоке. Однако его использование может создавать трудности, особенно в отношении безопасности отмены, что способно приводить к незаметным ошибкам, которые сложно отлаживать. Во многих случаях предпочтительнее альтернативы из экосистемы: они помогают избежать этих проблем благодаря более ясному синтаксису, предсказуемому управлению потоком выполнения и отсутствию необходимости вручную обрабатывать такие вопросы, как семантика слияния или безопасность отмены.
Объединение потоков
Если loop { select! { ... } } используется для опроса нескольких задач, объединение потоков представляет собой лаконичную альтернативу: оно изначально обеспечивает безопасную при отмене обработку и устраняет риск потери данных. Такие библиотеки, как tokio_stream, futures::stream и futures_concurrency, предоставляют средства для объединения потоков и последовательной обработки их результатов.
Пример с select!
struct File;
struct Channel;
struct Socket;
impl Socket {
async fn read_packet(&mut self) -> Vec<u8> {
vec![]
}
}
async fn read_send(_file: &mut File, _channel: &mut Channel) {
// do work that is not cancel safe
}
// open our IO types
let mut file = File;
let mut channel = Channel;
let mut socket = Socket;
loop {
tokio::select! {
_ = read_send(&mut file, &mut channel) => { /* ... */ },
_data = socket.read_packet() => { /* ... */ }
_ = futures::future::ready(()) => break
}
}Переход на merge
С помощью merge можно объединить несколько асинхронных задач в единый поток, избавившись от необходимости управлять задачами вручную и снизив риск нежелательного поведения, например потери данных.
use std::pin::pin;
use futures::stream::unfold;
use tokio_stream::StreamExt;
struct File;
struct Channel;
struct Socket;
impl Socket {
async fn read_packet(&mut self) -> Vec<u8> {
vec![]
}
}
async fn read_send(_file: &mut File, _channel: &mut Channel) {
// do work that is not cancel safe
}
enum Message {
Stop,
Sent,
Data(Vec<u8>),
}
// open our IO types
let file = File;
let channel = Channel;
let socket = Socket;
let a = unfold((file, channel), |(mut file, mut channel)| async {
read_send(&mut file, &mut channel).await;
Some((Message::Sent, (file, channel)))
});
let b = unfold(socket, |mut socket| async {
let data = socket.read_packet().await;
Some((Message::Data(data), socket))
});
let c = tokio_stream::iter([Message::Stop]);
let mut s = pin!(a.merge(b).merge(c));
while let Some(msg) = s.next().await {
match msg {
Message::Data(_data) => { /* ... */ }
Message::Sent => continue,
Message::Stop => break,
}
}Гонка будущих
Если нужно дождаться завершения первой из нескольких асинхронных задач, утилиты из экосистемы, такие как futures, futures-lite или futures-concurrency, предоставляют лаконичный синтаксис для организации гонки будущих:
futures_concurrency::future::Racefutures::select-
futures::stream::select_all(для потоков) futures_lite::future::orfutures_lite::future::race
use futures_concurrency::future::Race;
let task_a = async { Ok("ok") };
let task_b = async { Err("error") };
let result = (task_a, task_b).race().await;
match result {
Ok(output) => println!("First task completed with: {output}"),
Err(err) => eprintln!("Error occurred: {err}"),
}
MIT License
Copyright © Tokio Contributors
https://docs.rs/tokio/1.53.1/tokio/macro.select.html