Модуль sync
sync.Примитивы синхронизации для использования в асинхронных контекстах.
Программы на Tokio обычно организованы как набор задач, каждая из которых работает независимо и может выполняться в отдельном физическом потоке. Примитивы синхронизации, предоставляемые этим модулем, позволяют этим независимым задачам взаимодействовать друг с другом.
Передача сообщений
Наиболее распространённая форма синхронизации в программе на Tokio — передача сообщений. Две задачи работают независимо и отправляют сообщения друг другу для синхронизации. Это позволяет избежать общего состояния.
Передача сообщений реализована с помощью каналов. Канал поддерживает отправку сообщения от одной задачи-производителя одной или нескольким задачам-потребителям. Tokio предоставляет несколько видов каналов. Каждый вид канала поддерживает разные схемы передачи сообщений. Если канал поддерживает нескольких производителей, множество отдельных задач могут отправлять сообщения. Если канал поддерживает нескольких потребителей, множество отдельных задач могут получать сообщения.
Tokio предоставляет множество видов каналов, поскольку для разных схем передачи сообщений лучше подходят разные реализации.
oneshot-канал
oneshot-канал поддерживает отправку одного значения от одного производителя одному потребителю. Обычно этот канал используется для отправки результата вычисления ожидающей задаче.
Пример: использование oneshot-канала для получения результата вычисления.
use tokio::sync::oneshot;
async fn some_computation() -> String {
"represents the result of the computation".to_string()
}
let (tx, rx) = oneshot::channel();
tokio::spawn(async move {
let res = some_computation().await;
tx.send(res).unwrap();
});
// Do other work while the computation is happening in the background
// Wait for the computation result
let res = rx.await.unwrap();Обратите внимание: если задача выдаёт результат вычисления в качестве последнего действия перед завершением, для получения этого значения можно использовать JoinHandle вместо выделения ресурсов для oneshot-канала. Ожидание на JoinHandle возвращает Result. Если задача завершается с паникой, Joinhandle возвращает Err с причиной паники.
Пример:
async fn some_computation() -> String {
"the result of the computation".to_string()
}
let join_handle = tokio::spawn(async move {
some_computation().await
});
// Do other work while the computation is happening in the background
// Wait for the computation result
let res = join_handle.await.unwrap();
mpsc-канал
mpsc-канал поддерживает отправку множества значений от множества производителей одному потребителю. Этот канал часто используют для отправки работы задаче или для получения результатов множества вычислений.
Также следует использовать этот канал, если нужно отправлять множество сообщений от одного производителя одному потребителю. Специализированного канала spsc нет.
Пример: использование mpsc для потоковой передачи результатов последовательности вычислений.
use tokio::sync::mpsc;
async fn some_computation(input: u32) -> String {
format!("the result of computation {}", input)
}
let (tx, mut rx) = mpsc::channel(100);
tokio::spawn(async move {
for i in 0..10 {
let res = some_computation(i).await;
tx.send(res).await.unwrap();
}
});
while let Some(res) = rx.recv().await {
println!("got = {}", res);
}Аргумент mpsc::channel задаёт ёмкость канала. Это максимальное количество значений, которые могут одновременно храниться в канале в ожидании получения. Правильный выбор этого значения — ключ к созданию надёжных программ, поскольку ёмкость канала играет важную роль в обработке обратного давления.
Распространённая схема управления ресурсами при параллельном выполнении — запустить задачу, предназначенную для управления этим ресурсом, и использовать передачу сообщений между другими задачами для взаимодействия с ним. Ресурсом может быть всё, что нельзя использовать одновременно. Например, это может быть сокет или состояние программы. Так, если нескольким задачам нужно отправлять данные через один сокет, запустите задачу для управления сокетом и используйте канал для синхронизации.
Пример: отправка данных от множества задач через один сокет с помощью передачи сообщений.
use tokio::io::{self, AsyncWriteExt};
use tokio::net::TcpStream;
use tokio::sync::mpsc;
#[tokio::main]
async fn main() -> io::Result<()> {
let mut socket = TcpStream::connect("www.example.com:1234").await?;
let (tx, mut rx) = mpsc::channel(100);
for _ in 0..10 {
// Each task needs its own `tx` handle. This is done by cloning the
// original handle.
let tx = tx.clone();
tokio::spawn(async move {
tx.send(&b"data to write"[..]).await.unwrap();
});
}
// The `rx` half of the channel returns `None` once **all** `tx` clones
// drop. To ensure `None` is returned, drop the handle owned by the
// current task. If this `tx` handle is not dropped, there will always
// be a single outstanding `tx` handle.
drop(tx);
while let Some(res) = rx.recv().await {
socket.write_all(res).await?;
}
Ok(())
}Каналы mpsc и oneshot можно объединить, чтобы реализовать схему синхронизации «запрос/ответ» с общим ресурсом. Запускается задача для синхронизации ресурса, которая ожидает команды, получаемые через канал mpsc. Каждая команда содержит oneshot Sender, в который отправляется результат выполнения команды.
Пример: использование задачи для синхронизации счётчика u64. Каждая задача отправляет команду «получить и увеличить». Значение счётчика до увеличения отправляется через предоставленный oneshot-канал.
use tokio::sync::{oneshot, mpsc};
use Command::Increment;
enum Command {
Increment,
// Other commands can be added here
}
let (cmd_tx, mut cmd_rx) = mpsc::channel::<(Command, oneshot::Sender<u64>)>(100);
// Spawn a task to manage the counter
tokio::spawn(async move {
let mut counter: u64 = 0;
while let Some((cmd, response)) = cmd_rx.recv().await {
match cmd {
Increment => {
let prev = counter;
counter += 1;
response.send(prev).unwrap();
}
}
}
});
let mut join_handles = vec![];
// Spawn tasks that will send the increment command.
for _ in 0..10 {
let cmd_tx = cmd_tx.clone();
join_handles.push(tokio::spawn(async move {
let (resp_tx, resp_rx) = oneshot::channel();
cmd_tx.send((Increment, resp_tx)).await.ok().unwrap();
let res = resp_rx.await.unwrap();
println!("previous value = {}", res);
}));
}
// Wait for all tasks to complete
for join_handle in join_handles.drain(..) {
join_handle.await.unwrap();
}
broadcast-канал
broadcast-канал поддерживает отправку множества значений от множества производителей множеству потребителей. Каждый потребитель получит каждое значение. Этот канал можно использовать для реализации схемы «веерной рассылки», распространённой в системах pub/sub и чатах.
Этот канал используется реже, чем oneshot и mpsc, но у него всё же есть свои области применения.
Также следует использовать этот канал, если нужно передавать значения от одного производителя множеству потребителей. Специализированного канала широковещательной рассылки spmc нет.
Базовое использование
use tokio::sync::broadcast;
let (tx, mut rx1) = broadcast::channel(16);
let mut rx2 = tx.subscribe();
tokio::spawn(async move {
assert_eq!(rx1.recv().await.unwrap(), 10);
assert_eq!(rx1.recv().await.unwrap(), 20);
});
tokio::spawn(async move {
assert_eq!(rx2.recv().await.unwrap(), 10);
assert_eq!(rx2.recv().await.unwrap(), 20);
});
tx.send(10).unwrap();
tx.send(20).unwrap();
watch-канал
watch-канал поддерживает отправку множества значений от множества производителей множеству потребителей. Однако в канале хранится только самое последнее значение. Потребители получают уведомление при отправке нового значения, но нет гарантии, что они увидят все значения.
watch-канал похож на broadcast-канал ёмкостью 1.
Канал watch можно использовать для широковещательной рассылки изменений конфигурации или сигнализации об изменениях состояния программы, например о переходе к завершению работы.
Пример: использование watch-канала для уведомления задач об изменениях конфигурации. В этом примере файл конфигурации периодически проверяется. При изменении файла потребителям отправляется уведомление об изменениях конфигурации.
use tokio::sync::watch;
use tokio::time::{self, Duration, Instant};
use std::io;
#[derive(Debug, Clone, Eq, PartialEq)]
struct Config {
timeout: Duration,
}
impl Config {
async fn load_from_file() -> io::Result<Config> {
// file loading and deserialization logic here
}
}
async fn my_async_operation() {
// Do something here
}
// Load initial configuration value
let mut config = Config::load_from_file().await.unwrap();
// Create the watch channel, initialized with the loaded configuration
let (tx, rx) = watch::channel(config.clone());
// Spawn a task to monitor the file.
tokio::spawn(async move {
loop {
// Wait 10 seconds between checks
time::sleep(Duration::from_secs(10)).await;
// Load the configuration file
let new_config = Config::load_from_file().await.unwrap();
// If the configuration changed, send the new config value
// on the watch channel.
if new_config != config {
tx.send(new_config.clone()).unwrap();
config = new_config;
}
}
});
let mut handles = vec![];
// Spawn tasks that runs the async operation for at most `timeout`. If
// the timeout elapses, restart the operation.
//
// The task simultaneously watches the `Config` for changes. When the
// timeout duration changes, the timeout is updated without restarting
// the in-flight operation.
for _ in 0..5 {
// Clone a config watch handle for use in this task
let mut rx = rx.clone();
let handle = tokio::spawn(async move {
// Start the initial operation and pin the future to the stack.
// Pinning to the stack is required to resume the operation
// across multiple calls to `select!`
let op = my_async_operation();
tokio::pin!(op);
// Get the initial config value
let mut conf = rx.borrow().clone();
let mut op_start = Instant::now();
let sleep = time::sleep_until(op_start + conf.timeout);
tokio::pin!(sleep);
loop {
tokio::select! {
_ = &mut sleep => {
// The operation elapsed. Restart it
op.set(my_async_operation());
// Track the new start time
op_start = Instant::now();
// Restart the timeout
sleep.set(time::sleep_until(op_start + conf.timeout));
}
_ = rx.changed() => {
conf = rx.borrow_and_update().clone();
// The configuration has been updated. Update the
// `sleep` using the new `timeout` value.
sleep.as_mut().reset(op_start + conf.timeout);
}
_ = &mut op => {
// The operation completed!
return
}
}
}
});
handles.push(handle);
}
for handle in handles.drain(..) {
handle.await.unwrap();
}Синхронизация состояния
Остальные примитивы синхронизации предназначены для синхронизации состояния. Это асинхронные аналоги реализаций, предоставляемых std. Они работают так же, как соответствующие реализации std, но вместо блокировки потока ожидают асинхронно.
-
BarrierГарантирует, что несколько задач будут ждать друг друга, пока не достигнут определённой точки программы, после чего все продолжат выполнение одновременно. -
MutexМеханизм взаимного исключения, гарантирующий, что в каждый момент времени доступ к данным имеет не более одного потока. -
NotifyБазовое уведомление задачи.Notifyпозволяет уведомить ожидающую задачу без отправки данных. В этом случае задача пробуждается и возобновляет обработку. -
RwLockПредоставляет механизм взаимного исключения, который допускает одновременное чтение несколькими участниками, но разрешает запись только одному участнику за раз. В некоторых случаях это может быть эффективнее мьютекса. -
SemaphoreОграничивает степень параллелизма. Семафор хранит некоторое количество разрешений, которые задачи могут запрашивать для входа в критическую секцию. Семафоры полезны для реализации ограничений и лимитов любого рода.
Совместимость со средами выполнения
Все примитивы синхронизации, предоставляемые этим модулем, не зависят от среды выполнения. Их можно свободно перемещать между разными экземплярами среды выполнения Tokio и даже использовать в средах выполнения, отличных от Tokio.
При использовании в среде выполнения Tokio примитивы синхронизации участвуют в кооперативном планировании, чтобы избежать голодания. Эта возможность не применяется при использовании в средах выполнения, отличных от Tokio.
Исключение составляют методы, имена которых оканчиваются на _timeout: они зависят от среды выполнения, поскольку требуют доступа к таймеру Tokio. Дополнительную информацию об использовании каждого такого метода *_timeout см. в его документации.
Модули
- broadcast
- Широковещательная очередь с множеством производителей и потребителей. Каждое отправленное значение получают все потребители.
- futures
- Именованные типы future.
- mpsc
- Очередь с множеством производителей и одним потребителем для отправки значений между асинхронными задачами.
- oneshot
- Одноразовый канал используется для отправки одного сообщения между асинхронными задачами. Функция
channelиспользуется для создания пары дескрипторовSenderиReceiver, образующих канал. - watch
- Канал с множеством производителей и потребителей, в котором сохраняется только последнее отправленное значение.
Структуры
- Acquire
Error - Ошибка, возвращаемая функцией
Semaphore::acquire. - Barrier
- Барьер позволяет нескольким задачам синхронизировать начало вычисления.
- Barrier
Wait Result BarrierWaitResultвозвращается функциейwait, когда все задачи вBarrierвстретились на барьере.- Mapped
Mutex Guard - Дескриптор захваченного
Mutex, к которому была применена функция с помощьюMutexGuard::map. - Mutex
- Асинхронный тип, подобный
Mutex. - Mutex
Guard - Дескриптор захваченного
Mutex. Его можно удерживать во время любой точки.await, поскольку он реализуетSend. - Notify
- Уведомляет одну задачу о необходимости пробуждения.
- Once
Cell - Потокобезопасная ячейка, в которую можно записать значение только один раз.
- Owned
Mapped Mutex Guard - Владеющий дескриптор захваченного
Mutex, к которому была применена функция с помощьюOwnedMutexGuard::map. - Owned
Mutex Guard - Владеющий дескриптор захваченного
Mutex. - Owned
RwLock Mapped Write Guard - Владеющая структура RAII, освобождающая исключительный доступ к записи при удалении.
- Owned
RwLock Read Guard - Владеющая структура RAII, освобождающая общий доступ для чтения при удалении.
- Owned
RwLock Write Guard - Владеющая структура RAII, освобождающая исключительный доступ к записи при удалении.
- Owned
Semaphore Permit - Владеющее разрешение, полученное от семафора.
- RwLock
- Асинхронная блокировка для чтения и записи.
- RwLock
Mapped Write Guard - Структура RAII, освобождающая исключительный доступ к записи при удалении.
- RwLock
Read Guard - Структура RAII, освобождающая общий доступ для чтения при удалении.
- RwLock
Write Guard - Структура RAII, освобождающая исключительный доступ к записи при удалении.
- Semaphore
- Счётный семафор для асинхронного получения разрешений.
- Semaphore
Permit - Разрешение, полученное от семафора.
- SetOnce
- Потокобезопасная ячейка, в которую можно записать значение только один раз.
- SetOnce
Error - Ошибка, которая может быть возвращена функцией
SetOnce::set. - TryLock
Error - Ошибка, возвращаемая функциями
Mutex::try_lock,RwLock::try_readиRwLock::try_write.
Перечисления
- SetError
- Ошибки, которые могут быть возвращены функцией
OnceCell::set. - TryAcquire
Error - Ошибка, возвращаемая функцией
Semaphore::try_acquire.
MIT License
Copyright © Tokio Contributors
https://docs.rs/tokio/1.53.1/tokio/sync/index.html