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

实现TCP客户端随机等待并重连未监听服务器的问题

Tokio TCP客户端:自动重试连接与TCP保活配置问题

问题背景

我有一段基于Tokio的TCP客户端代码,需要实现服务器未监听时等待随机时长后重试连接的逻辑,同时之前配置的TCP Keepalive似乎没有生效。尝试了两种重试逻辑后,程序能打印连接成功的日志,但写入操作完全停滞,程序卡住。

原代码

use tokio::net::TcpStream;
use tokio::io::AsyncWriteExt;
use std::error::Error;
use socket2;
use std::{ time};

#[tokio::main]
pub async fn match_tcp_client(address: String, self_ip: String) -> Result<(), Box<dyn Error>> {
    // Connect to a peer
    println!("trying to connect from {} to address {}", self_ip, address);

    let mut stream = TcpStream::connect(address.clone()).await?;

    println!("connected from {} to address {}", self_ip, address);

    let sock_ref = socket2::SockRef::from(&stream);

    let mut ka = socket2::TcpKeepalive::new();
    ka = ka.with_time(time::Duration::from_secs(20));
    ka = ka.with_interval(time::Duration::from_secs(20));

    sock_ref.set_tcp_keepalive(&ka)?;

    // Write some data.
    stream.write_all(self_ip.as_bytes()).await?;
    stream.write_all(b"hello world!EOF").await?;
    
    Ok(())
}

尝试的重试代码(均出现写入停滞)

版本1

use tokio::net::TcpStream;
use tokio::io::AsyncWriteExt;
use std::error::Error;
use socket2;
use std::{ time};
use tokio::time::{ sleep, Duration};

#[tokio::main]
pub async fn match_tcp_client(address: String, self_ip: String) -> Result<(), Box<dyn Error>> {
    // Connect to a peer
    println!("trying to connect from {} to address {}", self_ip, address);

    while TcpStream::connect(address.clone()).await.is_err()
    {
        sleep(Duration::from_millis(10)).await;
    }

    let mut stream = TcpStream::connect(address.clone()).await?;

    println!("connected from {} to address {}", self_ip, address);

    let sock_ref = socket2::SockRef::from(&stream);

    let mut ka = socket2::TcpKeepalive::new();
    ka = ka.with_time(time::Duration::from_secs(20));
    ka = ka.with_interval(time::Duration::from_secs(20));

    sock_ref.set_tcp_keepalive(&ka)?;

    // Write some data.
    stream.write_all(self_ip.as_bytes()).await?;
    stream.write_all(b"hello world!EOF").await?;
    
    Ok(())
}

版本2

use tokio::net::TcpStream;
use tokio::io::AsyncWriteExt;
use std::error::Error;
use socket2;
use std::{ time};
use tokio::time::{ sleep, Duration};

#[tokio::main]
pub async fn match_tcp_client(address: String, self_ip: String) -> Result<(), Box<dyn Error>> {
    // Connect to a peer
    println!("trying to connect from {} to address {}", self_ip, address);

    loop{
        if TcpStream::connect(address.clone()).await.is_ok(){
            let mut stream = TcpStream::connect(address.clone()).await?;
            stream.set_linger(None)?;

            println!("connected from {} to address {}", self_ip, address);

            let sock_ref = socket2::SockRef::from(&stream);

            let mut ka = socket2::TcpKeepalive::new();
            ka = ka.with_time(time::Duration::from_secs(20));
            ka = ka.with_interval(time::Duration::from_secs(20));

            sock_ref.set_tcp_keepalive(&ka)?;

            // Write some data.
            stream.write_all(self_ip.as_bytes()).await?;
            stream.write_all(b"hello world!EOF").await?;
            
            stream.shutdown().await?;
            break;
        }
    }
    Ok(())
}

运行现象

在4台EC2实例上运行时,日志显示连接成功,但程序卡在写入环节,无任何进展。


问题分析

  1. 重复连接的资源浪费与阻塞风险:

    • 重试逻辑中,每次检查连接成功后又重新调用TcpStream::connect,不仅浪费资源,还可能因TCP半连接状态导致第二次连接阻塞。
    • 固定10ms的重试间隔过于密集,容易触发系统连接限制。
  2. TCP保活配置未生效:

    • 部分系统默认禁用TCP保活,需先启用再设置参数;socket2的API调用顺序需正确,部分平台对保活参数有额外限制。
  3. 写入停滞原因:

    • write_all会阻塞直到数据写入内核缓冲区,若服务器未读取数据导致缓冲区满,程序会无限阻塞。
    • 未处理连接后的读取操作,无法感知服务器的关闭信号,易陷入半开连接状态。

修正后的代码

use tokio::net::TcpStream;
use tokio::io::{AsyncWriteExt, AsyncReadExt};
use std::error::Error;
use socket2::{SockRef, TcpKeepalive};
use std::time;
use tokio::time::{sleep, Duration};
use rand::Rng;

#[tokio::main]
pub async fn match_tcp_client(address: String, self_ip: String) -> Result<(), Box<dyn Error>> {
    println!("trying to connect from {} to address {}", self_ip, &address);
    
    let mut rng = rand::thread_rng();
    // 重试逻辑:复用成功的连接,避免重复调用connect
    let mut stream = loop {
        match TcpStream::connect(&address).await {
            Ok(s) => {
                println!("connected from {} to address {}", self_ip, &address);
                break s;
            }
            Err(e) => {
                println!("connection failed: {}, retrying...", e);
                // 随机1-5秒重试,避免连接风暴
                let wait_ms = rng.gen_range(1000..5000);
                sleep(Duration::from_millis(wait_ms)).await;
            }
        }
    };

    // 修复TCP保活:先启用保活,再设置参数
    let sock_ref = SockRef::from(&stream);
    sock_ref.set_tcp_keepalive(true)?; // 确保保活功能开启
    
    let ka = TcpKeepalive::new()
        .with_time(time::Duration::from_secs(20))
        .with_interval(time::Duration::from_secs(20));
    
    sock_ref.set_tcp_keepalive(&ka)?;

    // 写入操作添加超时,防止无限阻塞
    tokio::select! {
        write_result = async {
            stream.write_all(self_ip.as_bytes()).await?;
            stream.write_all(b"hello world!EOF").await?;
            Ok(()) as Result<(), Box<dyn Error>>
        } => write_result?,
        _ = sleep(Duration::from_secs(30)) => {
            return Err("write operation timed out".into());
        }
    }

    // 读取服务器信号,处理连接关闭
    let mut buf = [0; 1024];
    match stream.read(&mut buf).await {
        Ok(n) if n == 0 => println!("server closed connection"),
        Ok(_) => println!("received response from server"),
        Err(e) => println!("read error: {}", e),
    }

    stream.shutdown().await?;
    Ok(())
}

关键改进点

  • 重试逻辑优化:

    • 直接复用成功的连接,避免重复调用connect
    • 采用1-5秒的随机重试间隔,防止多客户端同时重试引发的连接风暴
    • 打印连接错误信息,便于排查问题
  • TCP保活修复:

    • 先调用set_tcp_keepalive(true)启用保活功能,再设置具体参数,确保在所有平台生效
  • 写入停滞解决:

    • 使用tokio::select!为写入操作添加30秒超时,避免无限阻塞
    • 添加读取逻辑,感知服务器的关闭信号,避免半开连接
  • 其他细节:

    • 引入rand crate实现随机等待(需在Cargo.toml中添加rand = "0.8")
    • 减少字符串克隆,直接引用节省资源

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.22 23:52:02