实现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实例上运行时,日志显示连接成功,但程序卡在写入环节,无任何进展。
问题分析
重复连接的资源浪费与阻塞风险:
- 重试逻辑中,每次检查连接成功后又重新调用
TcpStream::connect,不仅浪费资源,还可能因TCP半连接状态导致第二次连接阻塞。 - 固定10ms的重试间隔过于密集,容易触发系统连接限制。
- 重试逻辑中,每次检查连接成功后又重新调用
TCP保活配置未生效:
- 部分系统默认禁用TCP保活,需先启用再设置参数;
socket2的API调用顺序需正确,部分平台对保活参数有额外限制。
- 部分系统默认禁用TCP保活,需先启用再设置参数;
写入停滞原因:
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秒超时,避免无限阻塞 - 添加读取逻辑,感知服务器的关闭信号,避免半开连接
- 使用
其他细节:
- 引入
randcrate实现随机等待(需在Cargo.toml中添加rand = "0.8") - 减少字符串克隆,直接引用节省资源
- 引入
内容的提问来源于stack exchange,提问作者Zubayr
相关产品推荐
相关产品推荐

