多EC2实例TCPListener线程通信异常:部分请求未被接收
分布式EC2实例TCP通信异常排查与修复
问题背景
我有4台EC2实例,计划构建分布式网络,要求每个实例向所有实例(含自身)发送数据。首先从文件读取IP地址到变量ip_address_clone,示例IP列表:
A.A.A.A B.B.B.B C.C.C.C D.D.D.D
随后通过线程为所有实例启动Server和Client,核心启动代码如下:
thread::scope(|s| { s.spawn(|| { for _ip in ip_address_clone.clone() { let _result = newserver::handle_server(INITIAL_PORT + port_count); } }); s.spawn(|| { let three_millis = time::Duration::from_millis(3); thread::sleep(three_millis); for ip in ip_address_clone.clone() { let self_ip_clone = self_ip.clone(); let _result = newclient::match_tcp_client( [ip.to_string(), (INITIAL_PORT + port_count).to_string()].join(":"), self_ip_clone, ); } }); });
Server代码
use std::error::Error; use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader}; use tokio::net::tcp::ReadHalf; use tokio::net::TcpListener; #[tokio::main] pub async fn handle_server(port: u32) -> Result<(), Box<dyn Error>> { let listener = TcpListener::bind(["0.0.0.0".to_string(), port.to_string()].join(":")) .await .unwrap(); // open connection let (mut socket, _) = listener.accept().await.unwrap(); // starts listening println!("---continue---"); let (reader, mut writer) = socket.split(); // tokio socket split to read and write concurrently let mut reader: BufReader<ReadHalf> = BufReader::new(reader); let mut line: String = String::new(); loop { //loop to get all the data from client until EOF is reached let _bytes_read: usize = reader.read_line(&mut line).await.unwrap(); if line.contains("EOF") //REACTOR to be used here { println!("EOF Reached"); writer.write_all(line.as_bytes()).await.unwrap(); println!("{}", line); line.clear(); break; } } Ok(()) }
Client代码
use std::error::Error; use tokio::io::AsyncWriteExt; use tokio::net::TcpStream; #[tokio::main] pub async fn match_tcp_client(address: String, self_ip: String) -> Result<(), Box<dyn Error>> { // Connect to a peer let mut stream = TcpStream::connect(address.clone()).await?; // Write some data. stream.write_all(self_ip.as_bytes()).await?; stream.write_all(b"hello world!EOF").await?; // stream.shutdown().await?; Ok(()) }
异常现象
实际运行时出现以下问题:
- 第一个启动的实例能接收所有请求
- 第二个实例无法接收第一个实例的请求
- 第三个实例无法接收前两个的请求,以此类推
第一个实例日志
Starting execution type nok launched ---continue--- EOF Reached A.A.A.Ahello world!EOF ---continue--- EOF Reached B.B.B.Bhello world!EOF ---continue--- EOF Reached C.C.C.Chello world!EOF ---continue--- EOF Reached D.D.D.Dhello world!EOF
第二个实例日志
Starting execution type nok launched ---continue--- EOF Reached B.B.B.Bhello world!EOF ---continue--- EOF Reached C.C.C.Chello world!EOF ---continue--- EOF Reached D.D.D.Dhello world!EOF
通信呈同步状态,每个实例仅能接收自身及后续IP实例的请求,Listener无法接收前置实例的连接请求。
问题根源
- Server单连接限制:
handle_server函数仅处理一个连接就结束,后续连接无法被监听。同时循环启动多个Server线程绑定同一个端口,后续线程的bind操作会失败(被unwrap()忽略),导致只有第一个Server线程在工作。 - 同步阻塞的启动逻辑:Server线程循环调用
handle_server,每次调用都要等一个连接处理完成才进入下一次循环,无法同时处理多个连接。 - 不可靠的Client启动时机:固定3ms的睡眠无法保证所有实例的Server都完成端口绑定,前置实例的Client可能在后置实例Server未就绪时发起连接,导致连接失败。
修复方案
1. 改造Server为多连接监听模式
让Server持续监听端口,为每个新连接启动独立异步任务处理:
use std::error::Error; use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader}; use tokio::net::TcpListener; pub async fn handle_server(port: u32) -> Result<(), Box<dyn Error>> { let listener = TcpListener::bind(format!("0.0.0.0:{}", port)).await?; println!("Server listening on port {}", port); // 持续监听新连接 loop { let (socket, addr) = listener.accept().await?; println!("---New connection from {}---", addr); // 为每个连接启动单独的异步任务 tokio::spawn(async move { let (reader, mut writer) = socket.split(); let mut reader = BufReader::new(reader); let mut line = String::new(); loop { match reader.read_line(&mut line).await { Ok(0) => { println!("Connection from {} closed", addr); break; } Ok(_) => { if line.contains("EOF") { println!("EOF Reached from {}", addr); writer.write_all(line.as_bytes()).await.unwrap(); println!("Received: {}", line); line.clear(); } } Err(e) => { eprintln!("Error reading from {}: {}", addr, e); break; } } } }); } }
2. 调整Server启动逻辑
每个实例只需启动一个Server线程监听端口即可,无需为每个IP重复启动:
thread::scope(|s| { // 启动单Server线程持续监听 s.spawn(|| { let rt = tokio::runtime::Runtime::new().unwrap(); rt.block_on(newserver::handle_server(INITIAL_PORT + port_count)).unwrap(); }); s.spawn(|| { let rt = tokio::runtime::Runtime::new().unwrap(); // 延长等待时间确保Server启动完成(或改用端口监听检测) thread::sleep(time::Duration::from_millis(100)); // 异步启动所有Client任务 for ip in ip_address_clone.clone() { let self_ip_clone = self_ip.clone(); let addr = format!("{}:{}", ip, INITIAL_PORT + port_count); rt.spawn(async move { let _ = newclient::match_tcp_client(addr, self_ip_clone).await; }); } }); });
3. 优化Client代码
添加连接超时,确保连接正常关闭:
use std::error::Error; use std::time::Duration; use tokio::io::AsyncWriteExt; use tokio::net::TcpStream; use tokio::time::timeout; pub async fn match_tcp_client(address: String, self_ip: String) -> Result<(), Box<dyn Error>> { // 设置5秒连接超时 let mut stream = timeout(Duration::from_secs(5), TcpStream::connect(&address)).await??; println!("Connected to {}", address); stream.write_all(self_ip.as_bytes()).await?; stream.write_all(b"hello world!EOF").await?; // 关闭写端告知Server数据发送完成 stream.shutdown().await?; Ok(()) }
4. 检查EC2网络配置
- 确保安全组允许
INITIAL_PORT + port_count端口的入站/出站流量(来源包含所有实例IP) - 确认实例间私有/公网IP可正常访问
内容的提问来源于stack exchange,提问作者Zubayr
相关产品推荐
相关产品推荐

