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

Rust异步中如何获取futures集合内最先完成的future结果

问题背景

需要同时驱动多个async_std::UdpSocket执行异步接收,任意一个socket收到数据时立刻返回对应接收结果,不需要等待所有socket都完成接收。
原有实现存在两个核心问题:

  • 错误使用futures::future::try_join_all:该方法会等待所有传入future全部执行完成后才返回结果集合,不符合“单任务完成即返回”的需求
  • 借用逻辑不合法:创建接收任务时提前从HashMap中获取缓冲区的可变引用,多个任务同时持有HashMap的可变引用会直接触发编译错误,且未将缓冲区所有权和对应socket绑定,无法安全跨异步任务传递

原有错误实现代码:

// build sockets and buffers
let sockets = Vec::new();
let buffers = HashMap::new();
for addr in ["127.0.0.1:4000", "10.0.0.1:8080", "192.168.0.1:9000"] {
    let socket = UdpSocket::bind(addr.parse().unwrap()).await?;
    let buf = [0u8; 1500];
    sockets.push(socket);
    buffers.insert(socket.peer_addr()?, buf);
}

// create an iterator of tasks reading from each socket
let socket_tasks = sockets.into_iter().map(|socket| {
    let socket_addr = socket.peer_addr().unwrap();
    let buf = buffers.get_mut(&socket_addr).unwrap();
    socket
        .recv_from(buf)
        .and_then(|_| async move { Ok(socket_addr) })
});

// wait for the first socket to return a value (DOESN'T WORK)
let buffer = try_join_all(socket_tasks)
    .and_then(|socket_addr| async { Ok(buffers.get(&socket_addr)) })
    .await
实现方法

直接使用futures::future::select_all即可满足需求,该方法的行为为:

  • 同时轮询驱动传入的所有future
  • 任意一个future完成时立刻终止等待,返回三元组:(完成future的返回值, 完成项在输入列表中的索引, 剩余未完成的future列表)
  • 不会阻塞等待其他未完成的future执行,和select!宏逻辑一致,但支持任意动态长度的future集合

实现时需要修正原有代码的借用问题:将socket、对应缓冲区、地址信息全部绑定到同一个异步块中,通过所有权转移避免跨任务借用冲突,不需要单独使用HashMap存储缓冲区。
可运行的实现代码:

use async_std::net::UdpSocket;
use futures::future::select_all;

#[async_std::main]
async fn main() -> std::io::Result<()> {
    let mut recv_futures = Vec::new();
    for addr in ["127.0.0.1:4000", "10.0.0.1:8080", "192.168.0.1:9000"] {
        let socket = UdpSocket::bind(addr.parse()?).await?;
        let bind_addr = socket.local_addr()?;
        let mut buf = [0u8; 1500];
        // 所有权全部转移进async块,不存在借用问题
        recv_futures.push(async move {
            let (recv_len, from_addr) = socket.recv_from(&mut buf).await?;
            Ok((bind_addr, from_addr, recv_len, buf))
        });
    }

    // 等待第一个完成的接收任务
    let (first_res, _idx, remaining_futures) = select_all(recv_futures).await;
    let (bind_addr, from_addr, recv_len, buf) = first_res?;
    
    println!("监听地址{bind_addr} 收到来自{from_addr}的{recv_len}字节数据");
    // 有效接收数据为缓冲区前recv_len字节
    let payload = &buf[..recv_len];

    // 如果需要继续等待其他socket的数据,保留remaining_futures即可,下次调用select_all传入即可
    Ok(())
}

注意事项

  • 未调用connect()的UDP监听socket不存在对端地址,不要调用peer_addr(),需要获取自身地址用local_addr()即可,获取发送方地址用recv_from返回的第二参数
  • 如果需要持续接收所有socket的数据包,不要丢弃select_all返回的剩余future列表,每次拿到一个结果后,将剩余列表传入下一次select_all调用即可实现持续的多socket数据包监听
  • 如果需要自动忽略接收出错的socket,直接等待第一个成功接收的结果,可以替换为futures::future::select_ok方法,不需要自己手动处理错误分支

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 23:12:17