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
相关产品推荐
相关产品推荐

