前言#
昨天我们稍微接触了下Stream,文档这部分应该也是还没完善的,说了几个方法就没了。
今天我们接着学
一次执行多个Future#
现在我们已经学会如何使用async/await了,也可以用于简单的异步场景了。
但是对于真实的异步情况来说一般需要并发的执行多个不同的异步运算。
接下来我们要接触一些方法涉及到上面的情况。
join!:等待所有futures都执行完毕。select!:等待其中一个future完成。Spawning:创建一个任务,它会携带一个future.(其实day2已经接触过了)FuturesUnordered:一组future,它们从它们子future收取(yields)结果
join#
futures::join这个宏的作用类似JoinHandle的join方法,可以等待并发的多个不同的异步future完成再结束。
join!宏#
先来看下例子
async fn get_book_and_music() -> (Book, Music) {
let book = get_book().await;
let music = get_music().await;
(book, music)
}这个异步函数返回了两个异步任务,但是这其实没有想象中的块,因为它俩都是异步函数,get_music会被get_book堵塞。
这个时候我们就可以使用futures::join!来将它俩并发处理。
use futures::join;
async fn get_book_and_music() -> (Book, Music) {
let book_fut = get_book();
let music_fut = get_music();
join!(book_fut, music_fut)
}虽然输出的时候还是要等待它俩完成才会输出,但是它俩在等待的过程是并发的,所以不会有堵塞。
来稍微看下源码

是一个声明宏,给传入的异步任务分别调用futures_macro::join_internal的宏。

这是个属性宏
真正处理的源码是下面这个
fn bind_futures(fut_exprs: Vec<Expr>, span: Span) -> (Vec<TokenStream2>, Vec<Ident>) {
let mut future_let_bindings = Vec::with_capacity(fut_exprs.len());
let future_names: Vec<_> = fut_exprs
.into_iter()
.enumerate()
.map(|(i, expr)| {
let name = format_ident!("_fut{}", i, span = span);
future_let_bindings.push(quote! {
// Move future into a local so that it is pinned in one place and
// is no longer accessible by the end user.
let mut #name = __futures_crate::future::maybe_done(#expr);
let mut #name = unsafe { __futures_crate::Pin::new_unchecked(&mut #name) };
});
name
})
.collect();
(future_let_bindings, future_names)
}
/// The `join!` macro.
pub(crate) fn join(input: TokenStream) -> TokenStream {
let parsed = syn::parse_macro_input!(input as Join);
// should be def_site, but that's unstable
let span = Span::call_site();
let (future_let_bindings, future_names) = bind_futures(parsed.fut_exprs, span);
let poll_futures = future_names.iter().map(|fut| {
quote! {
__all_done &= __futures_crate::future::Future::poll(
#fut.as_mut(), __cx).is_ready();
}
});
let take_outputs = future_names.iter().map(|fut| {
quote! {
#fut.as_mut().take_output().unwrap(),
}
});
TokenStream::from(quote! { {
#( #future_let_bindings )*
__futures_crate::future::poll_fn(move |__cx: &mut __futures_crate::task::Context<'_>| {
let mut __all_done = true;
#( #poll_futures )*
if __all_done {
__futures_crate::task::Poll::Ready((
#( #take_outputs )*
))
} else {
__futures_crate::task::Poll::Pending
}
}).await
} })
}简单的说这里将它俩合并成一个异步任务,然后等待它们完成。
原来的情况是music的会一直不动,直到第一个异步也就是book的ready之后才会执行。
现在变成它俩同时进入执行器的接收器队列中,这样就不会堵塞了。
注意这里手动调用了future的poll方法,所以外面也就不用调用await来处理了。
另外也不要这样做:其它语言中可以这样并发的做,但是在rust中不行,不会报错,但是依旧会堵塞。
// WRONG -- don't do this
async fn get_book_and_music() -> (Book, Music) {
let book_future = get_book();
let music_future = get_music();
(book_future.await, music_future.await)
}
try_join!#
如果你的future返回的数据类型是Result的,那么这里最好用try_join!来代替join!。
因为join!会一直执行下去,即使是这个数据返回的是Err变体。而try_join!会如果遇到Err变体会立即终止。
直接来看下用法
use futures::try_join;
async fn get_book() -> Result<Book, String> { /* ... */ Ok(Book) }
async fn get_music() -> Result<Music, String> { /* ... */ Ok(Music) }
async fn get_book_and_music() -> Result<(Book, Music), String> {
let book_fut = get_book();
let music_fut = get_music();
try_join!(book_fut, music_fut)
}注意future的Err类型得是一样的。
我们可以用futures::future::TryFutureExt里的.err_into或者..map_err(|e| ...)来合并错误类型
use futures::{
future::TryFutureExt,
try_join,
};
async fn get_book() -> Result<Book, ()> { /* ... */ Ok(Book) }
async fn get_music() -> Result<Music, String> { /* ... */ Ok(Music) }
async fn get_book_and_music() -> Result<(Book, Music), String> {
let book_fut = get_book().map_err(|()| "Unable to get book".to_string());
let music_fut = get_music();
try_join!(book_fut, music_fut)
} select!#
这个futures::select宏同样可以同时执行多个futures,不过不同于join!,select!只要有其中一个完成就会返回。
有些类似前端promise.all和promise.race的区别。
use futures::{
future::FutureExt, // for `.fuse()`
pin_mut,
select,
executor::block_on
};
async fn task_one() { /* ... */ }
async fn task_two() { /* ... */ }
async fn race_tasks() {
let t1 = task_one().fuse();
let t2 = task_two().fuse();
pin_mut!(t1, t2);
select! {
() = t1 => println!("task one completed first"),
() = t2 => println!("task two completed first"),
}
}
fn main () {
block_on(race_tasks());
}
pin_mut!`是一个声明宏相当于`Pin<&mut T>
select!`的语法是`<pattern> = <expression> => <code>default模式和complete模式用法#
select!支持default和complete两条分支:
default用在你模式啥都没匹配到但是已经结束咧~也就是都执行完了。complete用在当所有分支都完成并且都被选择了的场景。一般select!搭配循环且在结束的时候被调用。
use futures::{future, select, executor::{ block_on }};
async fn count() {
let mut a_fut = future::ready(4);
let mut b_fut = future::ready(6);
let mut total = 0;
loop {
select! {
a = a_fut => total += a,
b = b_fut => total += b,
complete => break,
default => unreachable!(), // never runs (futures are ready, then complete)
};
}
assert_eq!(total, 10);
}
fn main () {
block_on(count());
}搭配Unpin和FusedFuture#
上面第一个例子中,我们并没有提及.fuse这个方法。它用在原本应该是.await的地方。
它是必须调用的,因为使用select!的前提是future得实现Unpin和FusedFuture这俩trait。
实现Unpin的原因是select!的过程中并不是使用pin里包裹的值,而是一个可变引用,因为这样就不需要获取值的所有权,未完成的future就可以再次使用。
实现FusedFuture的原因类似Unpin,为了确保select在future执行结束之后不能再调用这个future的poll。FusedFuture用来跟踪future是否已经完成。上面第二个例子中future::ready返回的future实现了FusedFuture,那么这个时候就不会再poll了。(Stream中也是差不多的:FusedStream)
来看下例子
use futures::{
stream::{Stream, StreamExt, FusedStream},
select,
};
async fn add_two_streams(
mut s1: impl Stream<Item = u8> + FusedStream + Unpin,
mut s2: impl Stream<Item = u8> + FusedStream + Unpin,
) -> u8 {
let mut total = 0;
loop {
let item = select! {
x = s1.next() => x,
x = s2.next() => x,
complete => break,
};
if let Some(next_num) = item {
total += next_num;
}
}
total
}使用.next()/.try_next()将产生实现FusedFuture的future
在select!循环中使用Fused和Unordered来并发任务#
有一个很难描述但是很方便的函数叫Fuse::terminated(),它可以用来创建一个空的已经被终止的(terminated)future,这个空future将会被一个准备运行的future填充。
它能用于select循环过程中存储某个select自己产生的future。相当于select作用域外初始化一个future,等会在select循环过程中存储某个被创建的future。
我们先来看下Fuse::terminated的例子
use futures::{
future::{Fuse, FusedFuture, FutureExt},
stream::{FusedStream, Stream, StreamExt},
pin_mut,
select,
};
async fn get_new_num() -> u8 { /* ... */ 5 }
async fn run_on_new_num(_: u8) { /* ... */ }
async fn run_loop(
mut interval_timer: impl Stream<Item = ()> + FusedStream + Unpin,
starting_num: u8,
) {
let run_on_new_num_fut = run_on_new_num(starting_num).fuse();
let get_new_num_fut = Fuse::terminated();
pin_mut!(run_on_new_num_fut, get_new_num_fut);
loop {
select! {
() = interval_timer.select_next_some() => {
// The timer has elapsed. Start a new `get_new_num_fut`
// if one was not already running.
if get_new_num_fut.is_terminated() {
get_new_num_fut.set(get_new_num().fuse());
}
},
new_num = get_new_num_fut => {
// A new number has arrived -- start a new `run_on_new_num_fut`,
// dropping the old one.
run_on_new_num_fut.set(run_on_new_num(new_num).fuse());
},
// Run the `run_on_new_num_fut`
() = run_on_new_num_fut => {},
// panic if everything completed, since the `interval_timer` should
// keep yielding values indefinitely.
complete => panic!("`interval_timer` completed unexpectedly"),
}
}
}匹配到new_num的时候会被缓存到run_on_new_num_fut这个空future里。
当这个future完成后就会离开这个空future。
或者这个future倒计时结束后依旧还没运行,这个时候也会被踢出去。
可以调用is_terminated判断当前是什么状态。
注意.select_next_some()方法只能用于select时branch是Some(_)模式的场景,None会直接忽略。
然后再来看个FuturesUnordered的例子
use futures::{
future::{Fuse, FusedFuture, FutureExt},
stream::{FusedStream, FuturesUnordered, Stream, StreamExt},
pin_mut,
select,
};
async fn get_new_num() -> u8 { /* ... */ 5 }
async fn run_on_new_num(_: u8) -> u8 { /* ... */ 5 }
// Runs `run_on_new_num` with the latest number
// retrieved from `get_new_num`.
//
// `get_new_num` is re-run every time a timer elapses,
// immediately cancelling the currently running
// `run_on_new_num` and replacing it with the newly
// returned value.
async fn run_loop(
mut interval_timer: impl Stream<Item = ()> + FusedStream + Unpin,
starting_num: u8,
) {
let mut run_on_new_num_futs = FuturesUnordered::new();
run_on_new_num_futs.push(run_on_new_num(starting_num));
let get_new_num_fut = Fuse::terminated();
pin_mut!(get_new_num_fut);
loop {
select! {
() = interval_timer.select_next_some() => {
// The timer has elapsed. Start a new `get_new_num_fut`
// if one was not already running.
if get_new_num_fut.is_terminated() {
get_new_num_fut.set(get_new_num().fuse());
}
},
new_num = get_new_num_fut => {
// A new number has arrived -- start a new `run_on_new_num_fut`.
run_on_new_num_futs.push(run_on_new_num(new_num));
},
// Run the `run_on_new_num_futs` and check if any have completed
res = run_on_new_num_futs.select_next_some() => {
println!("run_on_new_num_fut returned {:?}", res);
},
// panic if everything completed, since the `interval_timer` should
// keep yielding values indefinitely.
complete => panic!("`interval_timer` completed unexpectedly"),
}
}
}和第一个例子不同,这里所有的future(相同的future)都会被执行并且会输出结果
TODO: spawning#
官方文档未完善
TODO: Cancellation and Timeouts#
同上
TODO: FuturesUnordered#
同上
总结#
莫得总结
这一章很潦草,即使是非TODO的,也就是简单的说一下怎么用。
参考#
- ^Executing-multiple-futures-at-a-time https://rust-lang.github.io/async-book/06_multiple_futures/01_chapter.html#executing-multiple-futures-at-a-time
- ^join https://rust-lang.github.io/async-book/06_multiple_futures/02_join.html#join
- ^join! https://rust-lang.github.io/async-book/06_multiple_futures/02_join.html#join-1
- ^try_join! https://rust-lang.github.io/async-book/06_multiple_futures/02_join.html#try_join
- ^select! https://rust-lang.github.io/async-book/06_multiple_futures/03_select.html#select
- ^https://rust-lang.github.io/async-book/06_multiple_futures/03_select.html#default---and-complete-- https://rust-lang.github.io/async-book/06_multiple_futures/03_select.html#default---and-complete--
- ^interation-with-Unpin-and-FusedFuture https://rust-lang.github.io/async-book/06_multiple_futures/03_select.html#interaction-with-unpin-and-fusedfuture
编辑于 2023-02-01 09:16・IP 属地广东
