use futures::stream::{self, StreamExt};
use std::future::Future;
use super::traits::AsyncIterator;
pub struct ParAsyncMap<I, F> {
iter: I,
concurrency: usize,
map_fn: F,
}
impl<I, F> ParAsyncMap<I, F> {
pub(super) fn new(iter: I, concurrency: usize, map_fn: F) -> Self {
Self {
iter,
concurrency,
map_fn,
}
}
}
impl<I, F, Fut, R> AsyncIterator for ParAsyncMap<I, F>
where
I: AsyncIterator,
F: Fn(I::Item) -> Fut + Send + Sync,
Fut: Future<Output = R> + Send,
R: Send,
{
type Item = R;
fn into_vec(self) -> Vec<Self::Item> {
let concurrency = self.concurrency.max(1);
let map_fn = self.map_fn;
let items = self.iter.into_vec();
futures::executor::block_on(async move {
stream::iter(items)
.map(|item| {
let map_fn = &map_fn;
async move { map_fn(item).await }
})
.buffered(concurrency)
.collect()
.await
})
}
}
pub struct ParAsyncFilter<I, F> {
iter: I,
concurrency: usize,
filter_fn: F,
}
impl<I, F> ParAsyncFilter<I, F> {
pub(super) fn new(iter: I, concurrency: usize, filter_fn: F) -> Self {
Self {
iter,
concurrency,
filter_fn,
}
}
}
impl<I, F, Fut> AsyncIterator for ParAsyncFilter<I, F>
where
I: AsyncIterator,
F: Fn(&I::Item) -> Fut + Send + Sync,
Fut: Future<Output = bool> + Send,
{
type Item = I::Item;
fn into_vec(self) -> Vec<Self::Item> {
let concurrency = self.concurrency.max(1);
let filter_fn = self.filter_fn;
let items = self.iter.into_vec();
futures::executor::block_on(async move {
stream::iter(items)
.map(|item| {
let filter_fn = &filter_fn;
async move {
let keep = filter_fn(&item).await;
(item, keep)
}
})
.buffered(concurrency)
.filter_map(|(item, keep)| async move { keep.then_some(item) })
.collect()
.await
})
}
}
pub(super) async fn for_each<I, F, Fut>(iter: I, concurrency: usize, func: F)
where
I: AsyncIterator,
F: Fn(I::Item) -> Fut + Send + Sync,
Fut: Future<Output = ()> + Send,
{
let items = iter.into_vec();
stream::iter(items)
.map(func)
.buffer_unordered(concurrency.max(1))
.for_each(|()| async {})
.await;
}