跳到主要内容

Module broadcast

搜索

Module broadcast 

Source
展开描述

一个多生产者、多消费者的广播队列。每个发送的值都会被所有消费者看到。

Sender 用于向所有连接的 Receiver 值广播值。 Sender 句柄是可克隆的, 允许并发的发送和接收操作。 SenderReceiver 都是 SendSync, 只要 TSend

当发送一个值时,所有Receiver 句柄都会被通知并接收该值。 该值在通道内仅存储一次, 并按需为每个 receiver 克隆。 一旦所有 receivers 都收到了该值的克隆, 该值将从通道中释放。

通过调用 channel 创建通道, 指定在任何给定时间通道可保留的最大消息数。

通过调用 Sender::subscribe 创建新的 Receiver 句柄。 返回的 Receiver 将接收调用 subscribe 之后发送的值。

此通道也适用于单生产者多消费者用例, 其中单个 sender 向多个 receivers 广播值。

§延迟(lagging)

由于发送的消息必须保留到所有Receiver 句柄都收到 一个克隆为止, 广播通道容易出现“慢接收者”问题。 在这种情况下, 除一个 receiver 之外的所有 receiver 都能以消息发送的速率接收值。 因为有一个 receiver 被卡住了, 通道开始积压。

此广播通道实现通过设置一个硬上限来解决此情况, 限制任何给定时间通道可保留的值的数量。 此上限作为参数传递给 channel 函数。

如果当通道达到容量时发送一个值, 通道当前持有的最旧值将被释放。 这为新值腾出了空间。 任何尚未看到已释放值的 receiver 在下次调用 recv 时 将返回 RecvError::Lagged

一旦返回 RecvError::Lagged, 滞后的 receiver 的位置 将被更新为通道中包含的最旧值。 下次调用 recv 将返回此值。

此行为使 receiver 能够检测到它何时已远远落后于 导致数据被丢弃。调用方可以决定如何响应: 通过中止其任务,或者通过容忍丢失的消息并 继续从通道消费。

§关闭(closing)

所有Sender 句柄都已被丢弃时, 将无法再发送新值。 此时,通道已“关闭”。 一旦 receiver 收到通道保留的所有值, 下次调用 recv 将返回 RecvError::Closed

Receiver 句柄被丢弃时, 任何尚未由该 receiver 读取的消息将被标记为已读。 如果该 receiver 是唯一 尚未读取该消息的 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;

let (tx, mut rx) = broadcast::channel(2);

tx.send(10).unwrap();
tx.send(20).unwrap();
tx.send(30).unwrap();

// The receiver lagged behind
assert!(rx.recv().await.is_err());

// At this point, we can abort or continue with lost messages

assert_eq!(20, rx.recv().await.unwrap());
assert_eq!(30, rx.recv().await.unwrap());

模块§

error
广播错误类型

结构体§

Receiver
broadcast channel 的接收端。
Sender
broadcast channel 的发送端。
WeakSender
一个 sender,不会阻止 channel 被关闭。

函数§

channel
创建一个有界的、多生产者、多消费者的 channel,每个发送的值都会广播给所有活跃的接收者。