4.2 生产级 Accept 循环

原文链接: https://book.async.rs/patterns/accept-loop.html

生产级的 accept 循环需要具备以下要素:

  1. 错误处理
  2. 限制并发连接数,以避免拒绝服务(DoS)攻击

错误处理

accept 循环中有两类错误:

  1. 单连接错误。系统用它们通知:队列中曾有连接,但已被对端丢弃。后续连接可能已在队列中,因此必须立即接受下一个连接。
  2. 资源不足。遇到此类错误时,立即接受下一个套接字没有意义。但监听器仍保持活跃,因此服务器应稍后再次尝试接受套接字。

以下是单连接错误的示例(正常模式与调试模式下的输出):

1
2
Error: Connection reset by peer (os error 104)
Error: Os { code: 104, kind: ConnectionReset, message: "Connection reset by peer" }

以下是最常见的资源不足错误示例:

1
2
Error: Too many open files (os error 24)
Error: Os { code: 24, kind: Other, message: "Too many open files" }

测试应用

要测试应用对这些错误的处理,可尝试以下步骤(仅适用于 Unix 系统)。

降低限制并启动应用:

1
2
3
4
5
6
$ ulimit -n 100
$ cargo run --example your_app
   Compiling your_app v0.1.0 (/work)
    Finished dev [unoptimized + debuginfo] target(s) in 5.47s
     Running `target/debug/examples/your_app`
Server is listening on: http://127.0.0.1:1234

然后在另一个终端运行 wrk 基准工具:

1
2
3
4
5
6
$ wrk -c 1000 http://127.0.0.1:1234
Running 10s test @ http://localhost:8080/
  2 threads and 1000 connections
$ telnet localhost 1234
Trying ::1...
Connected to localhost.

需要重点检查以下几点:

  1. 应用在出错时不应崩溃(但可以记录错误日志,见下文)
  2. 负载停止后(wrk 结束几秒后)应能再次连接到应用。上例中 telnet 即用于此目的,请确认输出包含 Connected to <hostname>。
  3. Too many open files 错误应记录到适当的日志中。为此,需将应用的「最大并发连接数」参数(见下文)设为大于 100 的值(针对本示例)。
  4. 测试期间检查应用的 CPU 占用。不应占满单个 CPU 核心的 100%(在 Rust 中,1000 个连接通常难以耗尽 CPU;若出现占满,说明错误处理不当)。

测试非 HTTP 应用

如有可能,使用相应的基准工具并设置合适的连接数。例如 redis-benchmark 提供 -c 参数,若你实现的是 Redis 协议即可使用。

或者,仍可使用 wrk,但需确保连接不会立即关闭。若会立即关闭,可在将连接交给协议处理器之前加入临时超时,例如:

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
16
17
18
19
20
# extern crate async_std;
# use std::time::Duration;
# use async_std::{
#     net::{TcpListener, ToSocketAddrs},
#     prelude::*,
# };
#
# type Result<T> = std::result::Result<T, Box<dyn std::error::Error + Send + Sync>>;
#
#async fn accept_loop(addr: impl ToSocketAddrs) -> Result<()> {
#    let listener = TcpListener::bind(addr).await?;
#    let mut incoming = listener.incoming();
while let Some(stream) = incoming.next().await {
    task::spawn(async {
        task::sleep(Duration::from_secs(10)).await; // 1
        connection_loop(stream).await;
    });
}
#     Ok(())
# }
  1. 确保 sleep 协程位于 spawn 的任务内部,而不是放在循环中。

手动处理错误

基本的 accept 循环可以写成这样:

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
# extern crate async_std;
# use std::time::Duration;
# use async_std::{
#     net::{TcpListener, ToSocketAddrs},
#     prelude::*,
# };
#
# type Result<T> = std::result::Result<T, Box<dyn std::error::Error + Send + Sync>>;
#
async fn accept_loop(addr: impl ToSocketAddrs) -> Result<()> {
    let listener = TcpListener::bind(addr).await?;
    let mut incoming = listener.incoming();
    while let Some(result) = incoming.next().await {
        let stream = match result {
            Err(ref e) if is_connection_error(e) => continue, // 1
            Err(e) => {
                eprintln!("Error: {}. Pausing for 500ms.", e); // 3
                task::sleep(Duration::from_millis(500)).await; // 2
                continue;
            }
            Ok(s) => s,
        };
        // body
    }
    Ok(())
}
  1. 忽略单连接错误。
  2. 资源不足时 sleep 后继续。
  3. 务必记录错误消息,因为这类错误通常表示系统配置不当,对运维人员很有帮助。

请务必测试你的应用。

外部 crate

async-listen crate 提供了实现此任务的辅助工具:

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
# extern crate async_std;
# extern crate async_listen;
# use std::time::Duration;
# use async_std::{
#     net::{TcpListener, ToSocketAddrs},
#     prelude::*,
# };
#
# type Result<T> = std::result::Result<T, Box<dyn std::error::Error + Send + Sync>>;
#
use async_listen::{ListenExt, error_hint};

async fn accept_loop(addr: impl ToSocketAddrs) -> Result<()> {

    let listener = TcpListener::bind(addr).await?;
    let mut incoming = listener
        .incoming()
        .log_warnings(log_accept_error) // 1
        .handle_errors(Duration::from_millis(500));
    while let Some(socket) = incoming.next().await { // 2
        // body
    }
    Ok(())
}

fn log_accept_error(e: &io::Error) {
    eprintln!("Error: {}. Listener paused for 0.5s. {}", e, error_hint(e)) // 3
}
  1. 记录资源不足(async-listen 称之为 warning)。若应用使用 log crate 或其他日志库,应写入日志。
  2. 经过 handle_errors 后,流产出的套接字不再包裹在 Result 中,因为所有错误已处理完毕。
  3. 除错误信息外,还会打印提示,向最终用户解释部分错误。例如,建议提高打开文件数限制并提供相关链接。

请务必测试你的应用。

连接数限制

即使已按错误处理一节完成所有措施,仍存在问题。

假设服务器需要打开文件来处理客户端请求。在某个时刻,可能出现以下情况:

  1. 客户端连接数已达到应用允许的最大文件描述符数量。
  2. 监听器收到 Too many open files 错误并进入 sleep。
  3. 某客户端通过先前已建立的连接发送请求。
  4. 为服务请求而打开文件失败,原因同样是 Too many open files 错误,直到其他客户端断开连接。

还有更多可能的情形;以上只是说明限制连接数非常有用的一例。一般而言,这是控制服务器所用资源、避免某些拒绝服务(DoS)攻击的方式之一。

async-listen crate

使用 async-listen 限制最大并发连接数的方式如下:

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
# extern crate async_std;
# extern crate async_listen;
# use std::time::Duration;
# use async_std::{
#     net::{TcpListener, TcpStream, ToSocketAddrs},
#     prelude::*,
# };
#
# type Result<T> = std::result::Result<T, Box<dyn std::error::Error + Send + Sync>>;
#
use async_listen::{ListenExt, Token, error_hint};

async fn accept_loop(addr: impl ToSocketAddrs) -> Result<()> {

    let listener = TcpListener::bind(addr).await?;
    let mut incoming = listener
        .incoming()
        .log_warnings(log_accept_error)
        .handle_errors(Duration::from_millis(500)) // 1
        .backpressure(100);
    while let Some((token, socket)) = incoming.next().await { // 2
         task::spawn(async move {
             connection_loop(&token, stream).await; // 3
         });
    }
    Ok(())
}
async fn connection_loop(_token: &Token, stream: TcpStream) { // 4
    // ...
}
# fn log_accept_error(e: &io::Error) {
#     eprintln!("Error: {}. Listener paused for 0.5s. {}", e, error_hint(e));
# }
  1. 需先处理错误,因为 backpressure 辅助函数期望的是 TcpStream 流,而非 Result。
  2. 新流产出的 token 由 backpressure 辅助函数计数。即,若丢弃 token,即可建立新连接。
  3. 将 token 的引用传给 connection loop,以将 token 的生命周期与连接的生命周期绑定。
  4. 函数中的 token 本身可忽略,故使用 _token。

请务必测试此行为。

最后修改 August 23, 2026: 更新 (499855b16)