目录
- 网络编程基础与标准库
- TCP 服务器与连接管理
- UDP 与无连接协议
- HTTP 客户端 reqwest
- gRPC 服务与 tonic
- WebSocket 实时通信
- Tokio 异步网络栈
- RESTful API 设计最佳实践
- 网络性能优化
- 速查表与工具链
1. 网络编程基础与标准库
Rust 标准库提供了跨平台的 socket 抽象,位于 std::net 模块:TcpListener、TcpStream、UdpSocket。它们是操作系统 socket 的薄封装,同步阻塞模型,适合作为网络编程的基础入门。
use std::net::TcpListener;
fn main() -> std::io::Result<()> {
let listener = TcpListener::bind("127.0.0.1:8080")?;
println!("监听端口 8080...");
for stream in listener.incoming() {
let stream = stream?;
println!("新连接: {:?}", stream.peer_addr()?);
// 每个连接创建一个线程处理(简单但开销大)
std::thread::spawn(|| handle(stream));
}
Ok(())
}
fn handle(mut stream: std::net::TcpStream) {
use std::io::{Read, Write};
let mut buf = [0u8; 1024];
if let Ok(n) = stream.read(&mut buf) {
let _ = stream.write_all(&buf[..n]); // 简单回声
}
}
同步 vs 异步:
| 模型 | 代表 | 并发能力 | 适用场景 |
|---|---|---|---|
| 线程阻塞 | std::net + 线程 | 每连接一线程,上千即告急 | 学习原型、低并发内部工具 |
| 多路复用 | mio/tokio | 单线程驱动上万连接 | 生产级高并发服务 |
| 异步运行时 | tokio/async-std | 零成本任务切换 | 现代 Rust 服务标配 |
结论:生产环境直接使用
tokio。std::net用于理解 socket 本质,但不要用它写高并发服务。
2. TCP 服务器与连接管理
TCP 是有连接、可靠、字节流的协议。服务端核心是 TcpListener::incoming() 循环 + 每连接的读写处理。
2.1 带优雅关闭的连接处理
use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader};
use tokio::net::{TcpListener, TcpStream};
use tokio::sync::broadcast;
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let listener = TcpListener::bind("0.0.0.0:8080").await?;
let (tx, _) = broadcast::channel::<String>(16); // 全局消息广播
loop {
let (stream, addr) = listener.accept().await?;
let tx = tx.clone();
tokio::spawn(async move {
let _ = handle_client(stream, addr, tx).await;
});
}
}
async fn handle_client(
mut stream: TcpStream,
addr: std::net::SocketAddr,
tx: broadcast::Sender<String>,
) -> Result<(), Box<dyn std::error::Error>> {
let (reader, mut writer) = stream.split();
let mut reader = BufReader::new(reader);
let mut lines = reader.lines();
// 读取与广播分离:读任务负责收消息,rx 负责下发
let mut rx = tx.subscribe();
loop {
tokio::select! {
line = lines.next_line() => {
match line? {
Some(l) => { let _ = tx.send(format!("{addr}: {l}")); }
None => break,
}
}
msg = rx.recv() => {
writer.write_all(msg?.as_bytes()).await?;
writer.write_all(b"\n").await?;
}
}
}
Ok(())
}
2.2 连接生命周期要点
| 要点 | 说明 |
|---|---|
accept() 循环 | 每个连接 spawn 一个任务,用 tokio::spawn |
拆读写作 split() | 读任务和写任务并发,避免死锁 |
| 心跳与超时 | tokio::time::timeout 包裹读写,防慢客户端拖垮资源 |
| 限流 | 全局并发连接数用 Semaphore 控制 |
| Graceful shutdown | ctrl_c() + broadcast 通知所有任务退出 |
3. UDP 与无连接协议
UDP 无连接、不可靠但低延迟,适合游戏、音视频、DNS 等场景。Rust 的 UdpSocket 一个 socket 即可收发。
use std::net::UdpSocket;
fn main() -> std::io::Result<()> {
let socket = UdpSocket::bind("0.0.0.0:9999")?;
let mut buf = [0u8; 1500];
loop {
let (n, src) = socket.recv_from(&mut buf)?;
println!("收到 {n} 字节来自 {src}");
socket.send_to(&buf[..n], src)?; // 回声
}
}
设计决策:
| 场景 | 推荐 | 原因 |
|---|---|---|
| 帧同步游戏 | UDP | 低延迟,丢包重发靠应用层 |
| 文件传输 | TCP | 可靠有序,丢包代价小 |
| 音视频 | UDP + RTP | 实时性优先,容忍少量丢包 |
| QUIC | UDP 之上 | HTTP/3,兼顾可靠与低延迟 |
4. HTTP 客户端 reqwest
reqwest 是事实标准的 HTTP 客户端,基于 hyper/tokio,支持异步、连接池、TLS、重定向、代理。
use reqwest::Client;
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
// 复用连接池:Client 要全局共享,不要每次新建
let client = Client::builder()
.timeout(std::time::Duration::from_secs(10))
.user_agent("my-app/1.0")
.pool_max_idle_per_host(20)
.build()?;
let resp = client
.get("https://api.example.com/users?page=1")
.header("Accept", "application/json")
.send()
.await?;
let status = resp.status();
let body: serde_json::Value = resp.json().await?; // 自动解析 JSON
println!("status={status}, name={}", body["data"]["name"]);
Ok(())
}
4.1 常见模式
// POST JSON
let resp = client
.post("https://api.example.com/users")
.json(&serde_json::json!({ "name": "Alice", "age": 25 }))
.send()
.await?;
// 查询参数
let url = reqwest::Url::parse_with_params(
"https://api.example.com/search",
&[("q", "rust"), ("page", "1")],
)?;
// 表单
let resp = client
.post("https://api.example.com/login")
.form(&[("username", "alice"), ("password", "secret")])
.send()
.await?;
4.2 错误处理分层
#[derive(Debug, thiserror::Error)]
enum ApiError {
#[error("网络错误: {0}")]
Network(#[from] reqwest::Error),
#[error("服务端返回 {0}: {1}")]
Server(reqwest::StatusCode, String),
}
async fn fetch_user(client: &Client, id: u64) -> Result<User, ApiError> {
let resp = client.get(format!("https://api.example.com/users/{id}")).send().await?;
match resp.status() {
reqwest::StatusCode::OK => Ok(resp.json().await?),
reqwest::StatusCode::NOT_FOUND => Err(ApiError::Server(resp.status(), "用户不存在".into())),
code => Err(ApiError::Server(code, resp.text().await.unwrap_or_default())),
}
}
5. gRPC 服务与 tonic
tonic 是 Rust 最流行的 gRPC 框架,基于 hyper + HTTP/2 + prost(Protocol Buffers)。
5.1 定义 proto
syntax = "proto3";
package helloworld;
service Greeter {
rpc SayHello (HelloRequest) returns (HelloReply);
rpc Chat (stream ChatMessage) returns (stream ChatMessage); // 双向流
}
message HelloRequest { string name = 1; }
message HelloReply { string message = 1; }
message ChatMessage { string text = 1; }
5.2 实现服务端
use tonic::{transport::Server, Request, Response, Status};
#[derive(Default)]
pub struct GreeterService;
#[tonic::async_trait]
impl helloworld::greeter_server::Greeter for GreeterService {
async fn say_hello(&self, req: Request<HelloRequest>) -> Result<Response<HelloReply>, Status> {
let name = req.into_inner().name;
Ok(Response::new(HelloReply {
message: format!("Hello, {name}!"),
}))
}
}
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let addr = "[::1]:50051".parse()?;
let svc = helloworld::greeter_server::GreeterServer::new(GreeterService::default());
Server::builder().add_service(svc).serve(addr).await?;
Ok(())
}
5.3 流式调用
// 双向流:接收 + 回显
async fn chat(
&self,
mut stream: Request<tonic::Streaming<ChatMessage>>,
) -> Result<Response<Self::ChatStream>, Status> {
let (tx, rx) = tokio::sync::mpsc::channel(8);
tokio::spawn(async move {
while let Some(msg) = stream.get_mut().message().await.unwrap_or(None) {
let _ = tx.send(Ok(ChatMessage { text: format!("echo: {}", msg.text) })).await;
}
});
Ok(Response::new(ReceiverStream::new(rx)))
}
gRPC vs REST 选型:
| 维度 | gRPC | REST |
|---|---|---|
| 协议 | HTTP/2 二进制 | HTTP/1.1 文本 |
| Schema | 强类型 .proto | OpenAPI |
| 流式 | 支持 4 种流 | 需 SSE/WebSocket 补 |
| 浏览器 | 需 gRPC-Web | 原生支持 |
| 场景 | 微服务内部 | 对外 API / 浏览器 |
6. WebSocket 实时通信
WebSocket 提供全双工长连接,适合聊天、通知、实时面板。Rust 常用 tokio-tungstenite(客户端+服务端)。
use tokio_tungstenite::tungstenite::Message;
use tokio_tungstenite::{accept_async, WebSocketStream};
use tokio::net::{TcpListener, TcpStream};
use futures_util::{SinkExt, StreamExt};
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let listener = TcpListener::bind("127.0.0.1:9001").await?;
while let Ok((stream, addr)) = listener.accept().await {
tokio::spawn(handle_ws(stream, addr));
}
Ok(())
}
async fn handle_ws(stream: TcpStream, addr: std::net::SocketAddr) {
let ws = accept_async(stream).await.unwrap();
println!("WS 连接: {addr}");
handle_messages(ws).await;
}
async fn handle_messages(mut ws: WebSocketStream<TcpStream>) {
while let Some(Ok(msg)) = ws.next().await {
match msg {
Message::Text(text) => {
let reply = format!("echo: {text}");
let _ = ws.send(Message::Text(reply.into())).await;
}
Message::Close(_) => break,
_ => {}
}
}
}
6.1 心跳保活
// 每 30 秒发送 ping,防止代理/负载均衡断开空闲连接
loop {
tokio::select! {
msg = ws.next() => {
if msg.is_none() { break; }
// 处理消息...
}
_ = tokio::time::sleep(Duration::from_secs(30)) => {
ws.send(Message::Ping(vec![].into())).await?;
}
}
}
7. Tokio 异步网络栈
Tokio 是异步运行时的事实标准,核心组件:
| 组件 | 职责 |
|---|---|
tokio::net | TcpListener/TcpStream/UdpSocket 异步封装 |
tokio::io | AsyncRead/AsyncWrite,AsyncBufReadExt 便捷方法 |
tokio::sync | 广播、mpsc、oneshot、Semaphore、Barrier |
tokio::select! | 多路等待事件,超时与取消 |
tokio::time | sleep、timeout、interval、速率限制 |
7.1 超时与取消
use tokio::time::timeout;
let result = timeout(Duration::from_secs(5), slow_operation()).await;
match result {
Ok(Ok(val)) => println!("成功: {val}"),
Ok(Err(e)) => println!("任务出错: {e}"),
Err(_) => println!("超时,已取消任务"),
}
7.2 连接池与背压
use tokio::sync::Semaphore;
struct Pool { sem: Semaphore }
impl Pool {
fn new(max: usize) -> Self { Self { sem: Semaphore::new(max) } }
async fn acquire(&self) -> Result<SemaphorePermit<'_>, ()> {
// 获取许可失败(无空位)即返回,实现背压
self.sem.try_acquire().map_err(|_| ())
}
}
背压原则:下游慢时,上游不要无限堆积。用有界队列(
mpsc::channel(n))+ Semaphore 限流,让系统自然限速而非 OOM。
8. RESTful API 设计最佳实践
结合 axum 实践一套生产级 API 设计规范:
use axum::{extract::{Path, State}, routing::{get, post}, Router, Json};
use std::sync::Arc;
#[derive(Clone)]
struct AppState {
users: Arc<dyn UserRepo + Send + Sync>,
}
#[derive(serde::Deserialize)]
struct CreateUser {
name: String,
#[serde(default)]
age: Option<u8>,
}
async fn create_user(
State(state): State<AppState>,
Json(input): Json<CreateUser>,
) -> Result<Json<User>, ApiError> {
let user = state.users.create(input).await?;
Ok(Json(user))
}
async fn get_user(
State(state): State<AppState>,
Path(id): Path<u64>,
) -> Result<Json<User>, ApiError> {
state.users.get(id).await?.ok_or(ApiError::NotFound)
}
pub fn router(state: AppState) -> Router {
Router::new()
.route("/users", post(create_user))
.route("/users/:id", get(get_user))
.with_state(state)
}
API 设计清单:
| 原则 | 做法 |
|---|---|
| 资源命名 | 复数名词 /users,层级 /users/{id}/orders |
| HTTP 动词 | GET 读 / POST 创建 / PUT 全量更新 / PATCH 部分更新 / DELETE 删除 |
| 状态码 | 200/201/204、400/401/403/404、429、500/502/503 |
| 分页 | ?limit=&cursor= 游标分页优于 offset |
| 幂等 | POST 创建带 Idempotency-Key |
| 版本 | /v1/ 前缀,破坏性变更升大版本 |
| 错误体 | 统一 { "code": "NOT_FOUND", "message": "..." } |
9. 网络性能优化
9.1 零拷贝 I/O
// 避免中间缓冲:直接从 socket 写入文件,绕过内核到用户的拷贝
use tokio::io::AsyncWriteExt;
async fn forward(mut src: tokio::net::TcpStream, dst: &mut tokio::fs::File) -> std::io::Result<()> {
// Rust 标准库/ tokio 通过 io_uring / splice 支持零拷贝
tokio::io::copy(&mut src, dst).await?;
Ok(())
}
9.2 实测指标(基准)
| 优化项 | 效果 |
|---|---|
| 连接复用(keep-alive) | RPS 提升 5-10x(省 TLS 握手 + 连接建立) |
| 连接池 | 高并发下 P99 降低 60%+ |
缓冲读写 BufReader | 小包吞吐提升 30%+ |
| 背压控制 | 系统稳定性:内存平稳,无 OOM |
loom/锁竞争降低 | 多核扩展性,无锁数据结构 |
# 压测
cargo install wrangler # 或其他压测工具
# 用 wrangler/http 压测
# 或:cargo install criterion 做微基准
10. 速查表与工具链
| 需求 | 推荐 crate |
|---|---|
| HTTP 服务端 | axum / actix-web / hyper |
| HTTP 客户端 | reqwest / hyper |
| gRPC | tonic + prost |
| WebSocket | tokio-tungstenite |
| 底层网络 | mio / tokio-net |
| DNS | hickory-dns(原 trust-dns) |
| 序列化 | serde + serde_json |
| 流量限速 | governor(令牌桶) |
| 压测 | wrangler-http / cargo-flamegraph |
一句话记忆:网络编程 = 选对模型(异步 Tokio)+ 管理连接(池化/背压/心跳)+ 设计好协议(REST/gRPC/WS),再辅以
select!处理超时与取消。
延伸阅读
- https://plumephp.com/rust-async-concurrency/ — Tokio 运行时与 Future 底层原理
- https://plumephp.com/rust-web-frameworks/ — Axum/Actix-web/Rocket 框架选型
- https://plumephp.com/posts/network/ — HTTP 缓存、反向代理与 API 网关
- https://plumephp.com/python-network-programming/ — Python 网络编程对比
- [[distributed-systems]] — 分布式系统中的网络通信与一致性
网络是分布式系统的地基。掌握 TCP/UDP 的本质、HTTP/gRPC/WS 协议语义、以及 Tokio 异步模型,是写出高并发可靠 Rust 服务的前提。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。