跳到主要内容

Module mpsc

搜索

Module mpsc 

Source
展开描述

用于在异步任务之间发送值的多生产者、单消费者队列。

本模块提供两种通道变体:有界和无界。 有界变体对通道可存储的消息数量有限制, 如果达到此限制, 尝试发送另一条消息将等待, 直到从通道接收到一条消息。 无界通道具有无限容量, 因此 send 方法将始终立即完成。 这使得 UnboundedSender 可以在同步和异步代码中使用。

类似于 std 提供的 mpsc 通道, 通道构造函数 为有界通道提供单独的发送和接收句柄 SenderReceiver, 为无界通道提供 UnboundedSenderUnboundedReceiver。 如果没有消息可读, 当发送新值时, 当前任务将收到通知。 SenderUnboundedSender 允许将值发送到通道中。 如果有界通道已达容量, 则 send 会被拒绝, 当有额外容量可用时, 任务将收到通知。 也就是说, 通道提供背压。

此通道也适用于单生产者单消费者用例。 (除非你只需要发送一条消息, 在这种情况下 应使用 oneshot 通道。)

§Disconnection

当所有 Sender 句柄 都已被丢弃时, 将无法再向通道发送值。 这被视为流的终止事件。 一旦所有 senders 都已被丢弃 且所有剩余的缓冲值 都已被接收, Receiver::recv 返回 NoneReceiver::poll_recv 返回 Poll::Ready(None))。

如果 Receiver 句柄 被丢弃, 则消息将无法再从通道中读取。 在这种情况下, 所有进一步的发送尝试 都将导致错误。 此外, 所有未读消息都将从通道中排出 并丢弃。

§Clean Shutdown

Receiver 被丢弃时, 未处理的消息 可能会保留在通道中。 通常, 更可取的做法是 执行“干净”的关闭。 为此, receiver 首先调用 close, 这将阻止任何进一步的消息 被发送到通道。 然后, receiver 消费完通道中的所有消息, 此时 receiver 可以被丢弃。

§Communicating between sync and async code

当你想在同步代码和异步代码之间通信时,需要考虑两种情况:

有界通道:如果你需要有界通道, 则应在两个通信方向上都使用 有界 Tokio mpsc 通道。 在同步代码中, 不是调用异步 sendrecv 方法, 而是需要使用 blocking_sendblocking_recv 方法。

无界通道:应使用与 receiver 所在位置 相匹配的通道类型。 因此, 若要从 async 向 sync 发送消息, 应使用 标准库的无界通道crossbeam。 类似地, 从 sync 向 async 发送消息, 应使用无界 Tokio mpsc 通道。

请注意, 上述说明是针对 mpsc 通道编写的, 但它们也可以推广到 其他类型的通道。 通常, 任何未标记为 async 的通道方法 都可以 在任何地方调用, 包括在运行时之外。 例如, 在运行时之外通过 oneshot 通道发送消息 是完全可行的。

§Multiple runtimes

mpsc 通道是与运行时无关的。 你可以自由地将它 在 Tokio 运行时的不同实例之间移动, 甚至可以从非 Tokio 运行时使用它。

在 Tokio 运行时中使用时, 它参与 协作式调度 以避免饥饿。 当从非 Tokio 运行时使用时,此功能不适用。

作为例外, 以 _timeout 结尾的方法不是与运行时无关的, 因为它们需要访问 Tokio 计时器。 有关其用法的更多信息, 请参阅每个 *_timeout 方法的文档。

§Allocation behavior

The implementation details described in this section may change in future Tokio releases.

mpsc 通道将元素存储在块中。块以链表形式组织。 发送将新元素推送到链表前端的块上, 接收则从后端的块上弹出它们。 在 64 位目标上一个块可以容纳 32 条消息, 在 32 位目标上一个块可以容纳 16 条消息。 这个数字与通道和消息大小无关。 每个块还存储 4 个指针大小的值用于簿记 (因此在 64 位机器上,每条消息有 1 字节的 开销)。

当块中的所有值都已被接收时,它变为空。 然后它将被释放,除非通道的第一个块 (正在存储新发送的元素的块)没有下一个块。 在这种情况下,空块将被重用为下一个块。

模块§

error
Channel 错误类型。

结构体§

OwnedPermit
owned 类型的 permit,用于向 channel 发送一个值。
Permit
用于向 channel 发送一个值的 permit。
PermitIterator
Iterator,对 Permit 进行迭代,可用于在 channel 中持有 n 个槽位。
Receiver
从关联的 Sender 接收值。
Sender
向关联的 Receiver 发送值。
UnboundedReceiver
从关联的 UnboundedSender 接收值。
UnboundedSender
向关联的 UnboundedReceiver 发送值。
WeakSender
一个 sender,不会阻止 channel 被关闭。
WeakUnboundedSender
一个无界 sender,不会阻止 channel 被关闭。

函数§

channel
创建一个有界 mpsc channel,用于在异步任务之间通过背压(backpressure)进行通信。
unbounded_channel
创建一个无界 mpsc channel,用于在异步任务之间通信(不提供背压)。