Модуль broadcast
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. - Weak
Sender - Отправитель, который не препятствует закрытию канала.
Функции
- channel
- Создаёт ограниченный канал с несколькими производителями и несколькими потребителями, в котором каждое отправленное значение передаётся всем активным получателям.
MIT License
Copyright © Tokio Contributors
https://docs.rs/tokio/1.53.1/tokio/sync/broadcast/index.html