Структура Receiver
pub struct Receiver<T> { /* private fields */ }
sync.Принимающая половина канала broadcast.
Не должна использоваться одновременно в нескольких местах. Сообщения можно получать с помощью recv.
Чтобы преобразовать этот приёмник в Stream, можно использовать обёртку BroadcastStream.
Примеры
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();Реализации
impl<T> Receiver<T>
pub fn len(&self) -> usize
Возвращает количество сообщений, отправленных в канал, которые этот Receiver ещё не получил.
В это число входят сообщения, которые уже были перезаписаны в кольцевом буфере и больше недоступны для чтения. Если len больше эффективной ёмкости канала (заданная ёмкость, округлённая вверх до ближайшей степени двойки), следующий вызов recv возвращает Err(RecvError::Lagged), а следующий вызов try_recv возвращает Err(TryRecvError::Lagged). Например, при channel(10) длина буфера равна 16, поэтому отставание начинается, когда len становится больше 16.
После успешного получения (в том числе после обработки Lagged и чтения сохранившихся сообщений) значение len соответствующим образом уменьшается.
Примеры
use tokio::sync::broadcast;
let (tx, mut rx1) = broadcast::channel(16);
tx.send(10).unwrap();
tx.send(20).unwrap();
assert_eq!(rx1.len(), 2);
assert_eq!(rx1.recv().await.unwrap(), 10);
assert_eq!(rx1.len(), 1);
assert_eq!(rx1.recv().await.unwrap(), 20);
assert_eq!(rx1.len(), 0);pub fn is_empty(&self) -> bool
Возвращает истину, если в канале нет сообщений, которые Receiver ещё не получил.
Примеры
use tokio::sync::broadcast;
let (tx, mut rx1) = broadcast::channel(16);
assert!(rx1.is_empty());
tx.send(10).unwrap();
tx.send(20).unwrap();
assert!(!rx1.is_empty());
assert_eq!(rx1.recv().await.unwrap(), 10);
assert_eq!(rx1.recv().await.unwrap(), 20);
assert!(rx1.is_empty());pub fn same_channel(&self, other: &Self) -> bool
Возвращает true, если приёмники принадлежат одному каналу.
Примеры
use tokio::sync::broadcast;
let (tx, rx) = broadcast::channel::<()>(16);
let rx2 = tx.subscribe();
assert!(rx.same_channel(&rx2));
let (_tx3, rx3) = broadcast::channel::<()>(16);
assert!(!rx3.same_channel(&rx2));pub fn sender_strong_count(&self) -> usize
Возвращает количество дескрипторов Sender.
pub fn sender_weak_count(&self) -> usize
Возвращает количество дескрипторов WeakSender.
impl<T: Clone> Receiver<T>
pub fn resubscribe(&self) -> Self
Повторно подписывается на канал, начиная с текущего последнего элемента.
Этот дескриптор Receiver получит копию всех значений, отправленных после повторной подписки. В него не войдут элементы, находящиеся в очереди текущего получателя. Рассмотрим следующий пример.
Примеры
use tokio::sync::broadcast;
let (tx, mut rx) = broadcast::channel(2);
tx.send(1).unwrap();
let mut rx2 = rx.resubscribe();
tx.send(2).unwrap();
assert_eq!(rx2.recv().await.unwrap(), 2);
assert_eq!(rx.recv().await.unwrap(), 1);pub async fn recv(&mut self) -> Result<T, RecvError>
Получает следующее значение для этого получателя.
Каждый дескриптор Receiver получит копию всех значений, отправленных после подписки.
Err(RecvError::Closed) возвращается, когда обе стороны Sender закрыты, что означает, что дальнейшая отправка значений в канал невозможна.
Если дескриптор Receiver отстаёт, то после заполнения канала новые отправленные значения перезаписывают старые в кольцевом буфере. Следующий вызов recv возвращает Err(RecvError::Lagged(n)), где n — это количество пропущенных получателем перезаписанных сообщений. Получатель остаётся подписанным; его внутренний курсор перемещается к самому старому значению, которое ещё хранится в канале. Последующий вызов recv возвращает это значение, если дальнейшие отправки не перезапишут его до того, как получатель его прочитает. Подробнее см. в разделе отставание.
Безопасность отмены
Этот метод безопасен при отмене. Если recv используется в качестве ветви в tokio::select!, а другая ветвь завершается первой, гарантируется, что в этом канале не было получено ни одного сообщения.
Примеры
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;
let (tx, mut rx) = broadcast::channel(2);
tx.send(10).unwrap();
tx.send(20).unwrap();
tx.send(30).unwrap();
// One message was overwritten before this receiver could read it.
assert!(matches!(rx.recv().await, Err(RecvError::Lagged(1))));
// Resume from the oldest retained message, or abort the task instead.
assert_eq!(20, rx.recv().await.unwrap());
assert_eq!(30, rx.recv().await.unwrap());pub fn try_recv(&mut self) -> Result<T, TryRecvError>
Пытается вернуть ожидающее значение для этого получателя, не дожидаясь его.
Это удобно для предварительной «оптимистичной проверки» перед тем, как решить, следует ли ожидать получение значения.
В отличие от recv, эта функция имеет три случая ошибки вместо двух (закрытый канал, пустой буфер и отстающий получатель).
Err(TryRecvError::Closed) возвращается, когда обе стороны Sender закрыты, что означает, что дальнейшая отправка значений в канал невозможна.
Если дескриптор Receiver отстаёт, то после заполнения канала новые отправленные значения перезаписывают старые в кольцевом буфере. Следующий вызов try_recv возвращает Err(TryRecvError::Lagged(n)), где n — это количество пропущенных получателем перезаписанных сообщений. Получатель остаётся подписанным; его внутренний курсор перемещается к самому старому значению, которое ещё хранится в канале. Последующий вызов try_recv возвращает это значение, если дальнейшие отправки не перезапишут его до того, как получатель его прочитает. Если для получения нет значений, возвращается Err(TryRecvError::Empty). Подробнее см. в разделе отставание.
Примеры
use tokio::sync::broadcast;
let (tx, mut rx) = broadcast::channel(16);
assert!(rx.try_recv().is_err());
tx.send(10).unwrap();
let value = rx.try_recv().unwrap();
assert_eq!(10, value);pub fn blocking_recv(&mut self) -> Result<T, RecvError>
Блокирующее получение для вызова вне асинхронного контекста.
Паника
Эта функция вызывает панику, если её вызвать внутри асинхронного контекста выполнения.
Примеры
use std::thread;
use tokio::sync::broadcast;
#[tokio::main]
async fn main() {
let (tx, mut rx) = broadcast::channel(16);
let sync_code = thread::spawn(move || {
assert_eq!(rx.blocking_recv(), Ok(10));
});
let _ = tx.send(10);
sync_code.join().unwrap();
}Реализации трейтов
Реализации авто-трейтов
impl<T> !RefUnwindSafe for Receiver<T>
impl<T> !UnwindSafe for Receiver<T>
impl<T> Freeze for Receiver<T>
impl<T> Send for Receiver<T>where T: Send,
impl<T> Sync for Receiver<T>where T: Send,
impl<T> Unpin for Receiver<T>
impl<T> UnsafeUnpin for Receiver<T>
Общие реализации
impl<T> BorrowMut<T> for Twhere T: ?Sized,
fn borrow_mut(&mut self) -> &mut T
impl<T> Instrument for T
fn instrument(self, span: Span) -> Instrumented<Self> ⓘ
fn in_current_span(self) -> Instrumented<Self> ⓘ
impl<T> WithSubscriber for T
fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self> ⓘ
fn with_current_subscriber(self) -> WithDispatch<Self> ⓘ
Subscriber к этому типу и возвращает обёртку WithDispatch. Подробнее
MIT License
Copyright © Tokio Contributors
https://docs.rs/tokio/1.53.1/tokio/sync/broadcast/struct.Receiver.html