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

多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无法接收前置实例的连接请求。

问题根源

  1. Server单连接限制:handle_server函数仅处理一个连接就结束,后续连接无法被监听。同时循环启动多个Server线程绑定同一个端口,后续线程的bind操作会失败(被unwrap()忽略),导致只有第一个Server线程在工作。
  2. 同步阻塞的启动逻辑:Server线程循环调用handle_server,每次调用都要等一个连接处理完成才进入下一次循环,无法同时处理多个连接。
  3. 不可靠的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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.22 23:32:03