Spec-Zone.ru › Tokio

Модуль broadcast

Доступно только при включённой функции crate sync.

Широковещательная очередь с несколькими производителями и несколькими потребителями. Каждое отправленное значение получают все потребители.

Sender используется для отправки значений всем подключённым получателям Receiver. Дескрипторы Sender можно клонировать, что позволяет выполнять отправку и получение одновременно. Sender и Receiver являются Send и Sync, пока T является Send.

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

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

Новые дескрипторы Receiver создаются вызовом Sender::subscribe. Возвращённый Receiver будет получать значения, отправленные после вызова subscribe.

Этот канал также подходит для сценария с одним производителем и несколькими потребителями, когда один отправитель передаёт значения многим получателям.

Отставание

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

Реализация этого широковещательного канала обрабатывает такой случай, устанавливая жёсткий верхний предел числа значений, которые канал может хранить в любой момент времени. Этот предел передаётся функции channel в качестве аргумента. Указанная ёмкость округляется вверх до ближайшей степени двойки; это округлённое значение определяет число сообщений, помещающихся в кольцевой буфер, и используется для обнаружения отставания. Например, channel(3) выделяет буфер длиной 4, поэтому получатель начинает отставать, только если отстаёт от отправителя более чем на 4 сообщения.

Если значение отправляется, когда канал заполнен, самое старое значение, хранящееся в канале, перезаписывается. Это освобождает место для нового значения. Любой получатель, который ещё не получил перезаписанное значение, при следующем вызове recv (или try_recv) вернёт RecvError::Lagged. Ошибка содержит число сообщений, отброшенных до позиции получателя, которые поэтому больше недоступны.

Возврат RecvError::Lagged не закрывает канал и не отключает получателя. Внутренняя позиция отстающего получателя перемещается к самому старому значению, которое ещё хранится в канале. При следующем успешном вызове recv / try_recv возвращается это самое старое сохранённое значение (если до его чтения получателем новые отправки не перезапишут его). Последующие вызовы получения продолжают возвращать значения в порядке отправки.

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

Закрытие

Когда все дескрипторы Sender удалены, отправка новых значений становится невозможной. В этот момент канал «закрывается». После того как получатель получит все значения, хранящиеся в канале, следующий вызов recv вернёт RecvError::Closed.

При удалении дескриптора Receiver все сообщения, не прочитанные получателем, помечаются как прочитанные. Если этот получатель был единственным, кто ещё не прочитал сообщение, оно в этот момент удаляется.

Примеры

Пример использования

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();

Обработка отставания

use tokio::sync::broadcast;
use tokio::sync::broadcast::error::RecvError;

// Capacity 2 → ring buffer of length 2.
let (tx, mut rx) = broadcast::channel(2);

tx.send(10).unwrap();
tx.send(20).unwrap();
// Overwrites 10; receiver has not read it yet.
tx.send(30).unwrap();

// One message (10) was dropped; cursor moves to the oldest retained value (20).
assert!(matches!(rx.recv().await, Err(RecvError::Lagged(1))));

// At this point, we can abort or continue with lost messages.
// Continuing resumes from the oldest retained message.
assert_eq!(20, rx.recv().await.unwrap());
assert_eq!(30, rx.recv().await.unwrap());

Модули

error
Типы ошибок широковещательного канала

Структуры

Receiver
Принимающая половина канала broadcast.
Sender
Отправляющая половина канала broadcast.
WeakSender
Отправитель, который не препятствует закрытию канала.

Функции

channel
Создаёт ограниченный канал с несколькими производителями и несколькими потребителями, в котором каждое отправленное значение передаётся всем активным получателям.

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

Spec-Zone.ru

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