3.4 发送消息

原文链接: https://book.async.rs/tutorial/sending_messages.html

发送消息

现在是时候实现另一半——发送消息了。 实现发送最明显的方法是让每个 connection_loop 访问其他每个客户端 TcpStream 的写半部分。 这样,客户端可以直接向接收者 .write_all 消息。 然而,这是错误的:如果 Alice 发送 bob: foo,而 Charley 发送 bob: bar,Bob 实际上可能收到 fobaor。 通过套接字发送消息可能需要多次系统调用,因此两个并发的 .write_all 可能会相互干扰!

经验法则是,每个 TcpStream 只应由单个任务写入。 所以让我们创建一个 connection_writer_loop 任务,它通过 channel 接收消息并将它们写入套接字。 该任务将是消息序列化的节点。 如果 Alice 和 Charley 同时向 Bob 发送两条消息,Bob 将看到与它们到达 channel 时相同顺序的消息。

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
# extern crate futures;
# use async_std::{
#     net::TcpStream,
#     prelude::*,
# };
use futures::channel::mpsc; // 1
use futures::sink::SinkExt;
use std::sync::Arc;

# type Result<T> = std::result::Result<T, Box<dyn std::error::Error + Send + Sync>>;
type Sender<T> = mpsc::UnboundedSender<T>; // 2
type Receiver<T> = mpsc::UnboundedReceiver<T>;

async fn connection_writer_loop(
    mut messages: Receiver<String>,
    stream: Arc<TcpStream>, // 3
) -> Result<()> {
    let mut stream = &*stream;
    while let Some(msg) = messages.next().await {
        stream.write_all(msg.as_bytes()).await?;
    }
    Ok(())
}
  1. 我们将使用 futures crate 中的 channel。
  2. 为简单起见,我们将使用无界 channel,本教程不会讨论背压。
  3. 由于 connection_loop 和 connection_writer_loop 共享同一个 TcpStream,我们需要将其放入 Arc。 注意,因为 client 只从流中读取,而 connection_writer_loop 只向流中写入,我们这里不会产生竞争。
最后修改 August 23, 2026: 更新 (499855b16)