You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用reqwest客户端与axum服务端时遭遇端口占用错误

端口占用错误(错误码10048)排查与解决

问题背景

  • 技术栈:reqwest客户端 + axum服务端,sqlx操作SQLite数据库
  • 客户端每秒发送约1000个请求,服务端单请求处理耗时约400ms
  • 此前为解决hyper::Error(IncompleteMessage),将reqwest客户端的pool_max_idle_per_host设为0
  • 运行15-45秒后随机触发错误码10048的端口占用问题,仅axum服务端占用3000端口,调整keep-alive参数仅改变崩溃时长

错误信息

Error: reqwest::Error {
    kind: Request, 
    url: Url { 
        scheme: "http", 
        cannot_be_a_base: false, 
        username: "", 
        password: None, 
        host: Some(Ipv4(127.0.0.1)), 
        port: Some(3000), 
        path: "/message", 
        query: None, 
        fragment: None 
    }, 
    source: hyper::Error(
        Connect, 
        ConnectError(
            "tcp connect error", 
            Os { 
                code: 10048, 
                kind: AddrInUse, 
                message: "Only one usage of each socket address (protocol/network address/port) is normally permitted." 
            }
        )
    )
}

相关代码片段

请求发送函数

// 发送速率约每秒1000次
pub async fn send_message_to_extractor(
    message: DiscordMessage,
    client: &reqwest::Client,
) -> Result<(), reqwest::Error> {
    let mut message = message;
    message.guild_id = Some("1234".to_owned());
    let request = client
        .post(EXTRACTOR_URL)
        .json(&message)
        .send()
        .await;
    match request {
        Ok(_) => Ok(()),
        Err(e) => Err(e),
    }
}

服务端处理函数

pub async fn handle_discord_message(pool: Extension<SqlitePool>, Json(payload): Json<DiscordMessage>) -> impl IntoResponse {
    ...
    // 解析消息约400ms耗时
    ...
    return StatusCode::OK;
}

客户端初始化及请求循环

// 调用send_message_to_extractor的代码,15-45秒后失败
let reqwest_client = reqwest::Client::builder()
    .pool_max_idle_per_host(0)
    .build()?;
while let Some(row) = stream.try_next().await? {
    let data: Vec<u8> = row.try_get("data")?;
    let data_string = if COMPRESSION {
        decode_reader(data)?
    } else {
        String::from_utf8(data)?
    };
    let discord_message: DiscordMessage = serde_json::from_str(&data_string)?;
    send_message_to_extractor(discord_message, &reqwest_client).await?;
    bar.inc(1);
}
Ok(())

服务端启动代码

// main.rs
let pool = SqlitePoolOptions::new()
    .max_connections(50)
    .connect(&DATABASE_URL)
    .await?;
let app = Router::new()
    .route("/message", post(routes::handle_discord_message))
    .layer(Extension(pool));
let addr = SocketAddr::from(([127, 0, 0, 1], 3000));
tracing::info!("listening on {}", addr);
axum::Server::bind(&addr)
    .serve(app.into_make_service())
    .await
    .unwrap();
Ok(())

问题原因

  1. 连接池禁用导致端口耗尽:pool_max_idle_per_host(0)彻底禁用TCP连接复用,每个请求都创建新连接。连接关闭后会进入TIME_WAIT状态(Windows默认240秒),短时间内大量TIME_WAIT套接字占用本地端口,最终耗尽所有可用临时端口,触发10048错误。
  2. 请求并发不匹配:客户端每秒发送1000请求,远超过服务端处理能力(单请求400ms,单核心每秒仅能处理2-3个请求),导致请求堆积,客户端连接等待时间变长,进一步加剧新连接的创建频率。
  3. SQLite性能瓶颈:SQLite是单文件数据库,即使设置50个连接,实际并发操作会被锁阻塞,放大服务端处理延迟,间接导致客户端连接积压。

解决方法

1. 恢复并优化连接池配置

取消pool_max_idle_per_host(0),启用连接池复用TCP连接,同时调整参数适配并发量:

use std::time::Duration;

let reqwest_client = reqwest::Client::builder()
    .pool_max_idle_per_host(100) // 根据并发需求调整,建议100-200
    .pool_idle_timeout(Duration::from_secs(60))
    .build()?;

若仍存在IncompleteMessage问题,补充TCP keep-alive参数:

let reqwest_client = reqwest::Client::builder()
    .pool_max_idle_per_host(100)
    .tcp_keepalive(Some(Duration::from_secs(30)))
    .build()?;

2. 控制客户端请求并发量

客户端当前串行发送请求,需用信号量限制并发数,匹配服务端处理能力:

use tokio::sync::Semaphore;
use std::sync::Arc;

// 设置并发数,建议等于服务端worker线程数+数据库连接数的合理值,比如50
let semaphore = Arc::new(Semaphore::new(50));
while let Some(row) = stream.try_next().await? {
    let data: Vec<u8> = row.try_get("data")?;
    let data_string = if COMPRESSION {
        decode_reader(data)?
    } else {
        String::from_utf8(data)?
    };
    let discord_message: DiscordMessage = serde_json::from_str(&data_string)?;
    
    let semaphore_clone = Arc::clone(&semaphore);
    let client_clone = reqwest_client.clone();
    tokio::spawn(async move {
        let _permit = semaphore_clone.acquire().await.unwrap();
        if let Err(e) = send_message_to_extractor(discord_message, &client_clone).await {
            eprintln!("请求发送失败: {}", e);
        }
    });
    
    bar.inc(1);
}
Ok(())

3. 优化服务端处理能力

  • 调整axum worker线程数:手动设置更多worker线程,利用多核资源:
axum::Server::bind(&addr)
    .workers(8) // 根据CPU核心数调整,建议核心数*2
    .serve(app.into_make_service())
    .await
    .unwrap();
  • 优化SQLite操作:
    • 启用WAL模式,提升并发性能:在数据库初始化时执行PRAGMA journal_mode=WAL;
    • 批量处理请求,减少数据库交互次数
    • 若高并发场景刚需,替换为PostgreSQL/MySQL等支持高并发的数据库

4. 调整Windows系统端口参数(可选)

若必须维持高并发连接,修改系统参数缩短TIME_WAIT超时并扩大临时端口范围:
以管理员身份打开命令提示符,执行:

# 将TIME_WAIT超时缩短至30秒
netsh int tcp set global timedwaitdelay=30
# 扩大TCP临时端口范围(从1024开始,共64511个端口)
netsh int ipv4 set dynamicport tcp start=1024 num=64511

内容的提问来源于stack exchange,提问作者Schaffn

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.22 09:37:56