4.2 生产级 Accept 循环
原文链接: https://book.async.rs/patterns/accept-loop.html
生产级的 accept 循环需要具备以下要素:
- 错误处理
- 限制并发连接数,以避免拒绝服务(DoS)攻击
错误处理
accept 循环中有两类错误:
- 单连接错误。系统用它们通知:队列中曾有连接,但已被对端丢弃。后续连接可能已在队列中,因此必须立即接受下一个连接。
- 资源不足。遇到此类错误时,立即接受下一个套接字没有意义。但监听器仍保持活跃,因此服务器应稍后再次尝试接受套接字。
以下是单连接错误的示例(正常模式与调试模式下的输出):
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.
|
需要重点检查以下几点:
- 应用在出错时不应崩溃(但可以记录错误日志,见下文)
- 负载停止后(
wrk 结束几秒后)应能再次连接到应用。上例中 telnet 即用于此目的,请确认输出包含 Connected to <hostname>。 Too many open files 错误应记录到适当的日志中。为此,需将应用的「最大并发连接数」参数(见下文)设为大于 100 的值(针对本示例)。- 测试期间检查应用的 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(())
# }
|
- 确保 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(())
}
|
- 忽略单连接错误。
- 资源不足时 sleep 后继续。
- 务必记录错误消息,因为这类错误通常表示系统配置不当,对运维人员很有帮助。
请务必测试你的应用。
外部 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
}
|
- 记录资源不足(
async-listen 称之为 warning)。若应用使用 log crate 或其他日志库,应写入日志。 - 经过
handle_errors 后,流产出的套接字不再包裹在 Result 中,因为所有错误已处理完毕。 - 除错误信息外,还会打印提示,向最终用户解释部分错误。例如,建议提高打开文件数限制并提供相关链接。
请务必测试你的应用。
连接数限制
即使已按错误处理一节完成所有措施,仍存在问题。
假设服务器需要打开文件来处理客户端请求。在某个时刻,可能出现以下情况:
- 客户端连接数已达到应用允许的最大文件描述符数量。
- 监听器收到
Too many open files 错误并进入 sleep。 - 某客户端通过先前已建立的连接发送请求。
- 为服务请求而打开文件失败,原因同样是
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));
# }
|
- 需先处理错误,因为
backpressure 辅助函数期望的是 TcpStream 流,而非 Result。 - 新流产出的 token 由 backpressure 辅助函数计数。即,若丢弃 token,即可建立新连接。
- 将 token 的引用传给 connection loop,以将 token 的生命周期与连接的生命周期绑定。
- 函数中的 token 本身可忽略,故使用
_token。
请务必测试此行为。