pub struct Sender<T> { /* private fields */ }展开描述
broadcast channel 的发送端。
可以从多个线程使用。
消息可以
通过
send
发送。
§示例
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();实现§
Source§impl<T> Sender<T>
impl<T> Sender<T>
Sourcepub fn send(&self, value: T) -> Result<usize, SendError<T>>
pub fn send(&self, value: T) -> Result<usize, SendError<T>>
尝试向所有活动的 Receiver 句柄发送一个值,如果无法发送则将其返回。
当至少有一个活动的 Receiver 句柄时,发送成功。发送失败的情况是所有关联的 Receiver 句柄都已被丢弃。
§Return
成功时,返回已订阅 Receiver 句柄的数量。这并不意味这么多接收者都会看到该消息,因为接收者可能在接收到消息之前就已被丢弃或落后(参见 lagging)。
§Note
返回 Ok 并不意味发送的值一定会被所有或任何活动的 Receiver 句柄观察到。Receiver 句柄可能在收到已发送消息之前就被丢弃。
返回 Err 并不意味未来的 send 调用会失败。可以通过调用 subscribe 来创建新的 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();Sourcepub fn subscribe(&self) -> Receiver<T>
pub fn subscribe(&self) -> Receiver<T>
创建一个新的 Receiver 句柄,它将接收在此次 subscribe 调用之后发送的值。
§示例
use tokio::sync::broadcast;
let (tx, _rx) = broadcast::channel(16);
// Will not be seen
tx.send(10).unwrap();
let mut rx = tx.subscribe();
tx.send(20).unwrap();
let value = rx.recv().await.unwrap();
assert_eq!(20, value);Sourcepub fn downgrade(&self) -> WeakSender<T>
pub fn downgrade(&self) -> WeakSender<T>
将 Sender 转换为 WeakSender,它不计入 RAII 语义——也就是说,如果通道的所有 Sender 实例都已被丢弃且仅剩 WeakSender 实例,则通道被关闭。
Sourcepub fn len(&self) -> usize
pub fn len(&self) -> usize
返回已排队值的数量。
一个值在排队中会一直保留,直到它被发送时所有存活的接收者都看到,或者被后续超出队列容量的发送操作从队列中驱逐。
§Note
与 Receiver::len 不同,此方法只报告已排队的值,而不报告在被所有接收者看到之前已被从队列中驱逐的值。
§示例
use tokio::sync::broadcast;
let (tx, mut rx1) = broadcast::channel(16);
let mut rx2 = tx.subscribe();
tx.send(10).unwrap();
tx.send(20).unwrap();
tx.send(30).unwrap();
assert_eq!(tx.len(), 3);
rx1.recv().await.unwrap();
// The len is still 3 since rx2 hasn't seen the first value yet.
assert_eq!(tx.len(), 3);
rx2.recv().await.unwrap();
assert_eq!(tx.len(), 2);Sourcepub fn is_empty(&self) -> bool
pub fn is_empty(&self) -> bool
如果没有排队的值则返回 true。
§示例
use tokio::sync::broadcast;
let (tx, mut rx1) = broadcast::channel(16);
let mut rx2 = tx.subscribe();
assert!(tx.is_empty());
tx.send(10).unwrap();
assert!(!tx.is_empty());
rx1.recv().await.unwrap();
// The queue is still not empty since rx2 hasn't seen the value.
assert!(!tx.is_empty());
rx2.recv().await.unwrap();
assert!(tx.is_empty());Sourcepub fn receiver_count(&self) -> usize
pub fn receiver_count(&self) -> usize
返回活动接收者的数量。
活动接收者是从 channel 或 subscribe 返回的 Receiver 句柄。这些句柄将接收此 Sender 发送的值。
§Note
不能保证已发送的消息一定会到达这么多接收者。活动接收者可能在丢弃之前再也不会调用 recv。
§示例
use tokio::sync::broadcast;
let (tx, _rx1) = broadcast::channel(16);
assert_eq!(1, tx.receiver_count());
let mut _rx2 = tx.subscribe();
assert_eq!(2, tx.receiver_count());
tx.send(10).unwrap();Sourcepub fn same_channel(&self, other: &Self) -> bool
pub fn same_channel(&self, other: &Self) -> bool
如果发送者属于同一通道则返回 true。
§示例
use tokio::sync::broadcast;
let (tx, _rx) = broadcast::channel::<()>(16);
let tx2 = tx.clone();
assert!(tx.same_channel(&tx2));
let (tx3, _rx3) = broadcast::channel::<()>(16);
assert!(!tx3.same_channel(&tx2));Sourcepub async fn closed(&self)
pub async fn closed(&self)
当订阅此 Sender 的 Receiver 数量达到零时完成的 future。
§示例
use futures::FutureExt;
use tokio::sync::broadcast;
let (tx, mut rx1) = broadcast::channel::<u32>(16);
let mut rx2 = tx.subscribe();
let _ = tx.send(10);
assert_eq!(rx1.recv().await.unwrap(), 10);
drop(rx1);
assert!(tx.closed().now_or_never().is_none());
assert_eq!(rx2.recv().await.unwrap(), 10);
drop(rx2);
assert!(tx.closed().now_or_never().is_some());Sourcepub fn strong_count(&self) -> usize
pub fn strong_count(&self) -> usize
返回 Sender 句柄的数量。
Sourcepub fn weak_count(&self) -> usize
pub fn weak_count(&self) -> usize
返回 WeakSender 句柄的数量。