Spec-Zone.ru › Tokio

Модуль sync

Доступно только при включённой функции crate feature 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
Канал с множеством производителей и потребителей, в котором сохраняется только последнее отправленное значение.

Структуры

AcquireError
Ошибка, возвращаемая функцией Semaphore::acquire.
Barrier
Барьер позволяет нескольким задачам синхронизировать начало вычисления.
BarrierWaitResult
BarrierWaitResult возвращается функцией wait, когда все задачи в Barrier встретились на барьере.
MappedMutexGuard
Дескриптор захваченного Mutex, к которому была применена функция с помощью MutexGuard::map.
Mutex
Асинхронный тип, подобный Mutex.
MutexGuard
Дескриптор захваченного Mutex. Его можно удерживать во время любой точки .await, поскольку он реализует Send.
Notify
Уведомляет одну задачу о необходимости пробуждения.
OnceCell
Потокобезопасная ячейка, в которую можно записать значение только один раз.
OwnedMappedMutexGuard
Владеющий дескриптор захваченного Mutex, к которому была применена функция с помощью OwnedMutexGuard::map.
OwnedMutexGuard
Владеющий дескриптор захваченного Mutex.
OwnedRwLockMappedWriteGuard
Владеющая структура RAII, освобождающая исключительный доступ к записи при удалении.
OwnedRwLockReadGuard
Владеющая структура RAII, освобождающая общий доступ для чтения при удалении.
OwnedRwLockWriteGuard
Владеющая структура RAII, освобождающая исключительный доступ к записи при удалении.
OwnedSemaphorePermit
Владеющее разрешение, полученное от семафора.
RwLock
Асинхронная блокировка для чтения и записи.
RwLockMappedWriteGuard
Структура RAII, освобождающая исключительный доступ к записи при удалении.
RwLockReadGuard
Структура RAII, освобождающая общий доступ для чтения при удалении.
RwLockWriteGuard
Структура RAII, освобождающая исключительный доступ к записи при удалении.
Semaphore
Счётный семафор для асинхронного получения разрешений.
SemaphorePermit
Разрешение, полученное от семафора.
SetOnce
Потокобезопасная ячейка, в которую можно записать значение только один раз.
SetOnceError
Ошибка, которая может быть возвращена функцией SetOnce::set.
TryLockError
Ошибка, возвращаемая функциями Mutex::try_lock, RwLock::try_read и RwLock::try_write.

Перечисления

SetError
Ошибки, которые могут быть возвращены функцией OnceCell::set.
TryAcquireError
Ошибка, возвращаемая функцией Semaphore::try_acquire.

MIT License
Copyright © Tokio Contributors
https://docs.rs/tokio/1.53.1/tokio/sync/index.html

Spec-Zone.ru

Настройки Оффлайн Что нового Помощь О нас
Spec-Zone .ru
спецификации, руководства, описания, API