Spec-Zone.ru › Tokio

Структура Semaphore

pub struct Semaphore { /* private fields */ }
Доступно только при включенной функции crate sync.

Счётный семафор, выполняющий асинхронное получение разрешений.

Семафор поддерживает набор разрешений. Разрешения используются для синхронизации доступа к общему ресурсу. Семафор отличается от мьютекса тем, что позволяет нескольким вызывающим сторонам одновременно обращаться к общему ресурсу.

Когда вызывается acquire и у семафора остаются разрешения, функция немедленно возвращает разрешение. Однако если свободных разрешений нет, acquire асинхронно ожидает, пока одно из занятых разрешений не будет освобождено. После этого освобождённое разрешение передаётся вызывающей стороне.

Этот Semaphore справедлив: разрешения выдаются в порядке поступления запросов. Это правило справедливости действует и при использовании acquire_many, поэтому, если вызов acquire_many в начале очереди запрашивает больше разрешений, чем доступно в данный момент, он может помешать завершению вызова acquire, даже если у семафора достаточно разрешений, чтобы завершить вызов acquire.

Чтобы использовать Semaphore в функции poll, можно воспользоваться утилитой PollSemaphore.

Порядок операций с памятью

Если задача записывает какие-либо данные, а затем освобождает разрешение, любая задача, которая впоследствии получит разрешение, гарантированно увидит эти данные. Благодаря этому семафор можно безопасно использовать для передачи данных между задачами через общее состояние.

Если выразить это точнее в терминах порядков памяти атомарных операций: получение разрешения (через acquire, acquire_many, try_acquire, try_acquire_many или их варианты _owned), освобождение разрешений (путём сброса SemaphorePermit или OwnedSemaphorePermit либо вызова add_permits или forget_permits) и закрытие семафора (через close) — это операции AcqRel. Они полностью упорядочены, и каждая из них синхронизируется со всеми такими операциями, которые предшествуют ей, обеспечивая те же гарантии, что и операции AcqRel над одним атомарным объектом.

Неудачная попытка получить разрешение (включая TryAcquireError::NoPermits и TryAcquireError::Closed), а также методы available_permits и is_closed ведут себя как загрузка Acquire.

Примеры

Простое использование:

use tokio::sync::{Semaphore, TryAcquireError};

let semaphore = Semaphore::new(3);

let a_permit = semaphore.acquire().await.unwrap();
let two_permits = semaphore.acquire_many(2).await.unwrap();

assert_eq!(semaphore.available_permits(), 0);

let permit_attempt = semaphore.try_acquire();
assert_eq!(permit_attempt.err(), Some(TryAcquireError::NoPermits));

Ограничение количества одновременно открытых файлов в программе

В большинстве операционных систем есть ограничения на количество открытых файловых дескрипторов. Даже в системах без явно заданных ограничений ресурсные ограничения косвенно устанавливают верхний предел числа открытых файлов. Если программа попытается открыть слишком много файлов и превысит этот предел, произойдёт ошибка.

В этом примере используется Semaphore с 100 разрешениями. Получая разрешение от Semaphore перед обращением к файлу, вы гарантируете, что программа будет открывать не более 100 файлов одновременно. При попытке открыть 101-й файл программа будет ждать, пока не появится свободное разрешение, и лишь затем продолжит открывать файл.

use std::io::Result;
use tokio::fs::File;
use tokio::sync::Semaphore;
use tokio::io::AsyncWriteExt;

static PERMITS: Semaphore = Semaphore::const_new(100);

async fn write_to_file(message: &[u8]) -> Result<()> {
    let _permit = PERMITS.acquire().await.unwrap();
    let mut buffer = File::create("example.txt").await?;
    buffer.write_all(message).await?;
    Ok(()) // Permit goes out of scope here, and is available again for acquisition
}

Ограничение количества одновременно отправляемых исходящих запросов

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

В этом примере используется Arc<Semaphore> с 10 разрешениями. Каждая созданная задача получает ссылку на семафор путём клонирования Arc<Semaphore>. Прежде чем отправить запрос, задача должна получить разрешение от семафора, вызвав Semaphore::acquire. Это гарантирует, что в любой момент времени параллельно отправляется не более 10 запросов. Отправив запрос, задача освобождает разрешение, чтобы другие задачи могли отправлять запросы.

use std::sync::Arc;
use tokio::sync::Semaphore;

// Define maximum number of parallel requests.
let semaphore = Arc::new(Semaphore::new(5));
// Spawn many tasks that will send requests.
let mut jhs = Vec::new();
for task_id in 0..50 {
    let semaphore = semaphore.clone();
    let jh = tokio::spawn(async move {
        // Acquire permit before sending request.
        let _permit = semaphore.acquire().await.unwrap();
        // Send the request.
        let response = send_request(task_id).await;
        // Drop the permit after the request has been sent.
        drop(_permit);
        // Handle response.
        // ...

        response
    });
    jhs.push(jh);
}
// Collect responses from tasks.
let mut responses = Vec::new();
for jh in jhs {
    let response = jh.await.unwrap();
    responses.push(response);
}
// Process responses.
// ...

Ограничение количества одновременно обрабатываемых входящих запросов

Как и количество одновременно открытых файлов, количество сетевых дескрипторов ограничено. Обработка неограниченного числа запросов может привести к отказу в обслуживании и многим другим проблемам.

В этом примере вместо глобальной переменной используется Arc<Semaphore>. Чтобы ограничить количество запросов, которые могут обрабатываться одновременно, мы получаем разрешение для каждой задачи до её создания. Получив разрешение, мы создаём новую задачу; после завершения задачи разрешение освобождается внутри неё, позволяя создавать другие задачи. Разрешения необходимо получать с помощью Semaphore::acquire_owned, чтобы их можно было перемещать между границами задач. (Поскольку наш семафор не является глобальной переменной — если бы он был глобальным, достаточно было бы acquire.)

use std::sync::Arc;
use tokio::sync::Semaphore;
use tokio::net::TcpListener;

#[tokio::main]
async fn main() -> std::io::Result<()> {
    let semaphore = Arc::new(Semaphore::new(3));
    let listener = TcpListener::bind("127.0.0.1:8080").await?;

    loop {
        // Acquire permit before accepting the next socket.
        //
        // We use `acquire_owned` so that we can move `permit` into
        // other tasks.
        let permit = semaphore.clone().acquire_owned().await.unwrap();
        let (mut socket, _) = listener.accept().await?;

        tokio::spawn(async move {
            // Do work using the socket.
            handle_connection(&mut socket).await;
            // Drop socket while the permit is still live.
            drop(socket);
            // Drop the permit, so more tasks can be created.
            drop(permit);
        });
    }
}

Запрет параллельного выполнения тестов

По умолчанию Rust запускает тесты из одного файла параллельно. Однако в некоторых случаях параллельное выполнение двух тестов может привести к проблемам. Например, такое может произойти, если тесты используют одну и ту же базу данных.

Рассмотрим следующий сценарий:

  1. test_insert: вставляет в базу данных пару «ключ-значение», а затем извлекает значение по тому же ключу, чтобы проверить вставку.
  2. test_update: вставляет ключ, затем обновляет его, присваивая новое значение, и проверяет, что значение было обновлено правильно.
  3. test_others: третий тест, который не изменяет состояние базы данных. Он может выполняться параллельно с другими тестами.

В этом примере test_insert и test_update должны выполняться последовательно, но порядок запуска тестов не имеет значения. Для решения этой задачи можно использовать семафор с одним разрешением.

use tokio::sync::Semaphore;

// Initialize a static semaphore with only one permit, which is used to
// prevent test_insert and test_update from running in parallel.
static PERMIT: Semaphore = Semaphore::const_new(1);

// Initialize the database that will be used by the subsequent tests.
static DB: Database = Database::setup();

#[tokio::test]
async fn test_insert() {
    // Acquire permit before proceeding. Since the semaphore has only one permit,
    // the test will wait if the permit is already acquired by other tests.
    let permit = PERMIT.acquire().await.unwrap();

    // Do the actual test stuff with database

    // Insert a key-value pair to database
    let (key, value) = ("name", 0);
    DB.insert(key, value).await;

    // Verify that the value has been inserted correctly.
    assert_eq!(DB.get(key).await, value);

    // Undo the insertion, so the database is empty at the end of the test.
    DB.delete(key).await;

    // Drop permit. This allows the other test to start running.
    drop(permit);
}

#[tokio::test]
async fn test_update() {
    // Acquire permit before proceeding. Since the semaphore has only one permit,
    // the test will wait if the permit is already acquired by other tests.
    let permit = PERMIT.acquire().await.unwrap();

    // Do the same insert.
    let (key, value) = ("name", 0);
    DB.insert(key, value).await;

    // Update the existing value with a new one.
    let new_value = 1;
    DB.update(key, new_value).await;

    // Verify that the value has been updated correctly.
    assert_eq!(DB.get(key).await, new_value);

    // Undo any modificattion.
    DB.delete(key).await;

    // Drop permit. This allows the other test to start running.
    drop(permit);
}

#[tokio::test]
async fn test_others() {
    // This test can run in parallel with test_insert and test_update,
    // so it does not use PERMIT.
}

Ограничение частоты с помощью корзины токенов

В этом примере демонстрируются методы add_permits и SemaphorePermit::forget.

Во многих приложениях и системах существуют ограничения на частоту выполнения определённых операций. Превышение этой частоты может привести к снижению производительности или даже к ошибкам.

В этом примере реализовано ограничение частоты с помощью корзины токенов. Корзина токенов — это способ ограничения частоты, который срабатывает не сразу, позволяя обрабатывать короткие всплески одновременно поступающих запросов.

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

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

Обратите внимание, что эта реализация неэффективна при малой длительности интервала, поскольку постоянно выполняет циклы и ожидание, расходуя много ресурсов CPU.

use std::sync::Arc;
use tokio::sync::Semaphore;
use tokio::time::{interval, Duration};

struct TokenBucket {
    sem: Arc<Semaphore>,
    jh: tokio::task::JoinHandle<()>,
}

impl TokenBucket {
    fn new(duration: Duration, capacity: usize) -> Self {
        let sem = Arc::new(Semaphore::new(capacity));

        // refills the tokens at the end of each interval
        let jh = tokio::spawn({
            let sem = sem.clone();
            let mut interval = interval(duration);
            interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);

            async move {
                loop {
                    interval.tick().await;

                    if sem.available_permits() < capacity {
                        sem.add_permits(1);
                    }
                }
            }
        });

        Self { jh, sem }
    }

    async fn acquire(&self) {
        // This can return an error if the semaphore is closed, but we
        // never close it, so this error can never happen.
        let permit = self.sem.acquire().await.unwrap();
        // To avoid releasing the permit back to the semaphore, we use
        // the `SemaphorePermit::forget` method.
        permit.forget();
    }
}

impl Drop for TokenBucket {
    fn drop(&mut self) {
        // Kill the background task so it stops taking up resources when we
        // don't need it anymore.
        self.jh.abort();
    }
}

let capacity = 5;
let update_interval = Duration::from_secs_f32(1.0 / capacity as f32);
let bucket = TokenBucket::new(update_interval, capacity);

for _ in 0..5 {
    bucket.acquire().await;

    // do the operation
}

Реализации

impl Semaphore

pub const MAX_PERMITS: usize = super::batch_semaphore::Semaphore::MAX_PERMITS

Максимальное количество разрешений, которое может хранить семафор. Это usize::MAX >> 3.

Превышение этого предела обычно приводит к панике.

pub fn new(permits: usize) -> Self

Создаёт новый семафор с начальным количеством разрешений.

Вызывает панику, если permits превышает Semaphore::MAX_PERMITS.

pub const fn const_new(permits: usize) -> Self

Создаёт новый семафор с начальным количеством разрешений.

При использовании нестабильной функции tracing unstable feature объект Semaphore, созданный с помощью const_new, не будет инструментирован. Поэтому он не будет виден в tokio-console. Если требуется инструментированный объект, для его создания следует использовать Semaphore::new.

Примеры
use tokio::sync::Semaphore;

static SEM: Semaphore = Semaphore::const_new(10);

pub fn available_permits(&self) -> usize

Возвращает текущее количество доступных разрешений.

pub fn add_permits(&self, n: usize)

Добавляет семафору n новых разрешений.

Максимальное количество разрешений — Semaphore::MAX_PERMITS; эта функция вызовет панику, если предел будет превышен.

pub fn forget_permits(&self, n: usize) -> usize

Уменьшает количество разрешений семафора не более чем на n.

Если разрешений недостаточно и уменьшить их количество на n невозможно, возвращает фактическое количество удалённых разрешений.

pub async fn acquire(&self) -> Result<SemaphorePermit<'_>, AcquireError>

Получает разрешение от семафора.

Если семафор закрыт, возвращает AcquireError. В противном случае возвращает SemaphorePermit, представляющий полученное разрешение.

Безопасность при отмене

Этот метод использует очередь, чтобы справедливо распределять разрешения в порядке поступления запросов. Отмена вызова acquire приводит к потере места в очереди.

Примеры
use tokio::sync::Semaphore;

let semaphore = Semaphore::new(2);

let permit_1 = semaphore.acquire().await.unwrap();
assert_eq!(semaphore.available_permits(), 1);

let permit_2 = semaphore.acquire().await.unwrap();
assert_eq!(semaphore.available_permits(), 0);

drop(permit_1);
assert_eq!(semaphore.available_permits(), 1);

pub async fn acquire_many( &self, n: u32, ) -> Result<SemaphorePermit<'_>, AcquireError>

Получает от семафора n разрешений.

Если семафор закрыт, возвращает AcquireError. В противном случае возвращает SemaphorePermit, представляющий полученные разрешения.

Безопасность при отмене

Этот метод использует очередь, чтобы справедливо распределять разрешения в порядке поступления запросов. Отмена вызова acquire_many приводит к потере места в очереди.

Примеры
use tokio::sync::Semaphore;

let semaphore = Semaphore::new(5);

let permit = semaphore.acquire_many(3).await.unwrap();
assert_eq!(semaphore.available_permits(), 2);

pub fn try_acquire(&self) -> Result<SemaphorePermit<'_>, TryAcquireError>

Пытается получить разрешение от семафора.

Если семафор был закрыт, этот метод возвращает TryAcquireError::Closed, а если разрешений не осталось — TryAcquireError::NoPermits. В противном случае возвращается SemaphorePermit, представляющее полученные разрешения.

Примеры
use tokio::sync::{Semaphore, TryAcquireError};

let semaphore = Semaphore::new(2);

let permit_1 = semaphore.try_acquire().unwrap();
assert_eq!(semaphore.available_permits(), 1);

let permit_2 = semaphore.try_acquire().unwrap();
assert_eq!(semaphore.available_permits(), 0);

let permit_3 = semaphore.try_acquire();
assert_eq!(permit_3.err(), Some(TryAcquireError::NoPermits));

pub fn try_acquire_many( &self, n: u32, ) -> Result<SemaphorePermit<'_>, TryAcquireError>

Пытается получить n разрешений от семафора.

Если семафор был закрыт, этот метод возвращает TryAcquireError::Closed, а если разрешений осталось недостаточно — TryAcquireError::NoPermits. В противном случае возвращается SemaphorePermit, представляющее полученные разрешения.

Примеры
use tokio::sync::{Semaphore, TryAcquireError};

let semaphore = Semaphore::new(4);

let permit_1 = semaphore.try_acquire_many(3).unwrap();
assert_eq!(semaphore.available_permits(), 1);

let permit_2 = semaphore.try_acquire_many(2);
assert_eq!(permit_2.err(), Some(TryAcquireError::NoPermits));

pub async fn acquire_owned( self: Arc<Self>, ) -> Result<OwnedSemaphorePermit, AcquireError>

Получает разрешение от семафора.

Чтобы вызвать этот метод, семафор необходимо обернуть в Arc. Если семафор был закрыт, метод возвращает AcquireError. В противном случае возвращается OwnedSemaphorePermit, представляющее полученное разрешение.

Безопасность отмены

Этот метод использует очередь, чтобы справедливо распределять разрешения в порядке поступления запросов. Отмена вызова acquire_owned приведёт к потере места в очереди.

Примеры
use std::sync::Arc;
use tokio::sync::Semaphore;

let semaphore = Arc::new(Semaphore::new(3));
let mut join_handles = Vec::new();

for _ in 0..5 {
    let permit = semaphore.clone().acquire_owned().await.unwrap();
    join_handles.push(tokio::spawn(async move {
        // perform task...
        // explicitly own `permit` in the task
        drop(permit);
    }));
}

for handle in join_handles {
    handle.await.unwrap();
}

pub async fn acquire_many_owned( self: Arc<Self>, n: u32, ) -> Result<OwnedSemaphorePermit, AcquireError>

Получает n разрешений от семафора.

Чтобы вызвать этот метод, семафор необходимо обернуть в Arc. Если семафор был закрыт, метод возвращает AcquireError. В противном случае возвращается OwnedSemaphorePermit, представляющее полученное разрешение.

Безопасность отмены

Этот метод использует очередь, чтобы справедливо распределять разрешения в порядке поступления запросов. Отмена вызова acquire_many_owned приведёт к потере места в очереди.

Примеры
use std::sync::Arc;
use tokio::sync::Semaphore;

let semaphore = Arc::new(Semaphore::new(10));
let mut join_handles = Vec::new();

for _ in 0..5 {
    let permit = semaphore.clone().acquire_many_owned(2).await.unwrap();
    join_handles.push(tokio::spawn(async move {
        // perform task...
        // explicitly own `permit` in the task
        drop(permit);
    }));
}

for handle in join_handles {
    handle.await.unwrap();
}

pub fn try_acquire_owned( self: Arc<Self>, ) -> Result<OwnedSemaphorePermit, TryAcquireError>

Пытается получить разрешение от семафора.

Чтобы вызвать этот метод, семафор необходимо обернуть в Arc. Если семафор закрыт, возвращается TryAcquireError::Closed; если разрешений не осталось, возвращается TryAcquireError::NoPermits. В противном случае возвращается OwnedSemaphorePermit, представляющее полученное разрешение.

Примеры
use std::sync::Arc;
use tokio::sync::{Semaphore, TryAcquireError};

let semaphore = Arc::new(Semaphore::new(2));

let permit_1 = Arc::clone(&semaphore).try_acquire_owned().unwrap();
assert_eq!(semaphore.available_permits(), 1);

let permit_2 = Arc::clone(&semaphore).try_acquire_owned().unwrap();
assert_eq!(semaphore.available_permits(), 0);

let permit_3 = semaphore.try_acquire_owned();
assert_eq!(permit_3.err(), Some(TryAcquireError::NoPermits));

pub fn try_acquire_many_owned( self: Arc<Self>, n: u32, ) -> Result<OwnedSemaphorePermit, TryAcquireError>

Пытается получить n разрешений от семафора.

Чтобы вызвать этот метод, семафор необходимо обернуть в Arc. Если семафор закрыт, возвращается TryAcquireError::Closed; если разрешений не осталось, возвращается TryAcquireError::NoPermits. В противном случае возвращается OwnedSemaphorePermit, представляющее полученное разрешение.

Примеры
use std::sync::Arc;
use tokio::sync::{Semaphore, TryAcquireError};

let semaphore = Arc::new(Semaphore::new(4));

let permit_1 = Arc::clone(&semaphore).try_acquire_many_owned(3).unwrap();
assert_eq!(semaphore.available_permits(), 1);

let permit_2 = semaphore.try_acquire_many_owned(2);
assert_eq!(permit_2.err(), Some(TryAcquireError::NoPermits));

pub fn close(&self)

Закрывает семафор.

Это не позволяет семафору выдавать новые разрешения и уведомляет всех ожидающих.

Примеры
use tokio::sync::Semaphore;
use std::sync::Arc;
use tokio::sync::TryAcquireError;

let semaphore = Arc::new(Semaphore::new(1));
let semaphore2 = semaphore.clone();

tokio::spawn(async move {
    let permit = semaphore.acquire_many(2).await;
    assert!(permit.is_err());
    println!("waiter received error");
});

println!("closing semaphore");
semaphore2.close();

// Cannot obtain more permits
assert_eq!(semaphore2.try_acquire().err(), Some(TryAcquireError::Closed))

pub fn is_closed(&self) -> bool

Возвращает true, если семафор закрыт

Реализации трейтов

impl Debug for Semaphore

fn fmt(&self, f: &mut Formatter<'_>) -> Result

Форматирует значение с помощью заданного форматировщика. Подробнее

Автоматические реализации трейтов

impl !Freeze for Semaphore

impl !RefUnwindSafe for Semaphore

impl !UnwindSafe for Semaphore

impl Send for Semaphore

impl Sync for Semaphore

impl Unpin for Semaphore

impl UnsafeUnpin for Semaphore

Общие реализации

impl<T> Any for T
where T: 'static + ?Sized,

fn type_id(&self) -> TypeId

Получает TypeId значения self. Подробнее

impl<T> Borrow<T> for T
where T: ?Sized,

fn borrow(&self) -> &T

Неизменяемо заимствует значение, находящееся во владении. Подробнее

impl<T> BorrowMut<T> for T
where T: ?Sized,

fn borrow_mut(&mut self) -> &mut T

Изменяемо заимствует значение, находящееся во владении. Подробнее

impl<T> From<T> for T

fn from(t: T) -> T

Возвращает аргумент без изменений.

impl<T> Instrument for T

fn instrument(self, span: Span) -> Instrumented<Self> ⓘ

Инструментирует этот тип с помощью предоставленного Span, возвращая оболочку Instrumented. Подробнее

fn in_current_span(self) -> Instrumented<Self> ⓘ

Инструментирует этот тип с помощью текущего Span, возвращая оболочку Instrumented. Подробнее

impl<T, U> Into<U> for T
where U: From<T>,

fn into(self) -> U

Вызывает U::from(self).

То есть эта преобразование выполняет действия, выбранные реализацией From<T> for U.

impl<T, U> TryFrom<U> for T
where U: Into<T>,

type Error = Infallible

Тип, возвращаемый в случае ошибки преобразования.

fn try_from(value: U) -> Result<T, <T as TryFrom<U>>::Error>

Выполняет преобразование.

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

type Error = <U as TryFrom<T>>::Error

Тип, возвращаемый в случае ошибки преобразования.

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Выполняет преобразование.

impl<T> WithSubscriber for T

fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self> ⓘ
where S: Into<Dispatch>,

Присоединяет указанный Subscriber к этому типу и возвращает обёртку WithDispatch. Подробнее

fn with_current_subscriber(self) -> WithDispatch<Self> ⓘ

Присоединяет к этому типу текущий по умолчанию Subscriber и возвращает обёртку WithDispatch. Подробнее

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

Spec-Zone.ru

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