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

tokio::select是否支持任意数量分支?能否基于Vec<TcpListener>实现?

基于Vec实现多端口监听的Tokio Echo服务器

当然可以实现基于Vec<TcpListener>的动态多端口监听,不需要为每个端口单独定义变量。下面提供两种实用方案:

方案一:用通道统一接收连接(推荐)

这种方式逻辑简单,扩展性强,适合任意数量的端口。我们为每个listener启动独立的异步任务,循环调用accept,并将接收到的连接通过通道发送给主循环处理。

use tokio::{io, net, spawn, sync::mpsc};
use std::sync::Arc;

#[tokio::main]
async fn main() {
    // 定义要监听的端口列表
    let ports = vec!["127.0.0.1:8001", "127.0.0.1:8002"];
    let listeners: Vec<_> = ports.into_iter()
        .map(|port| Arc::new(net::TcpListener::bind(port).await.unwrap()))
        .collect();

    // 创建通道,用于传递接收到的连接
    let (tx, mut rx) = mpsc::channel(10);

    // 为每个listener启动独立监听任务
    for listener in listeners {
        let tx = tx.clone();
        spawn(async move {
            loop {
                let (conn, _) = listener.accept().await.unwrap();
                tx.send(conn).await.unwrap();
            }
        });
    }

    // 主循环处理所有连接
    while let Some(conn) = rx.recv().await {
        spawn(handle(conn));
    }
}

async fn handle(mut conn: net::TcpStream) {
    let (mut read, mut write) = conn.split();
    io::copy(&mut read, &mut write).await.unwrap();
}

说明

  • 用Arc包裹TcpListener,确保每个异步任务都能安全共享监听实例。
  • 每个listener在独立任务中循环accept,不会互相阻塞。
  • 主循环只需监听通道,处理所有端口的连接,逻辑清晰。

方案二:用select_all动态选择完成的任务

如果你更贴近原生select!的逻辑,可以使用futures库提供的select_all函数,它能从一组future中选出第一个完成的任务。

首先需要在Cargo.toml中添加依赖:

[dependencies]
tokio = { version = "1.0", features = ["full"] }
futures = "0.3"

然后实现代码:

use tokio::{io, net, spawn};
use futures::future::select_all;
use std::sync::Arc;

#[tokio::main]
async fn main() {
    let ports = vec!["127.0.0.1:8001", "127.0.0.1:8002"];
    let listeners: Vec<_> = ports.into_iter()
        .map(|port| Arc::new(net::TcpListener::bind(port).await.unwrap()))
        .collect();

    // 初始化所有listener的accept future
    let mut futures: Vec<_> = listeners.iter()
        .map(|listener| listener.accept())
        .collect();

    loop {
        if futures.is_empty() {
            break;
        }

        // 等待第一个完成的accept操作
        let (result, idx, mut remaining_futures) = select_all(futures).await;
        let (conn, _) = result.unwrap();
        
        // 启动任务处理连接
        spawn(handle(conn));
        
        // 将当前listener的下一个accept future重新加入列表,继续监听
        let listener = &listeners[idx];
        remaining_futures.push(listener.accept());
        
        // 更新future列表,进入下一轮循环
        futures = remaining_futures;
    }
}

async fn handle(mut conn: net::TcpStream) {
    let (mut read, mut write) = conn.split();
    io::copy(&mut read, &mut write).await.unwrap();
}

说明

  • select_all返回第一个完成的结果、对应的索引,以及剩余未完成的future。
  • 每次处理完一个连接后,需要将该listener的下一个accept future重新加入列表,保证端口持续监听。
  • 同样用Arc共享TcpListener,避免所有权问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 13:10:25