moirai_iter/async_iter/
parallel.rs1use futures::stream::{self, StreamExt};
4use std::future::Future;
5
6use super::traits::AsyncIterator;
7use crate::stream::{retained_buffered, retained_unordered};
8
9pub struct ParAsyncMap<I, F> {
11 iter: I,
12 concurrency: usize,
13 map_fn: F,
14}
15
16impl<I, F> ParAsyncMap<I, F> {
17 pub(super) fn new(iter: I, concurrency: usize, map_fn: F) -> Self {
18 Self {
19 iter,
20 concurrency,
21 map_fn,
22 }
23 }
24}
25
26impl<I, F, Fut, R> AsyncIterator for ParAsyncMap<I, F>
27where
28 I: AsyncIterator,
29 F: Fn(I::Item) -> Fut + Send + Sync,
30 Fut: Future<Output = R> + Send,
31 R: Send,
32{
33 type Item = R;
34
35 fn into_vec(self) -> Vec<Self::Item> {
36 let concurrency = self.concurrency.max(1);
37 let map_fn = self.map_fn;
38 let items = self.iter.into_vec();
39 futures::executor::block_on(async move {
40 let futures = stream::iter(items).map(|item| {
41 let map_fn = &map_fn;
42 async move { map_fn(item).await }
43 });
44 retained_buffered(futures, concurrency).collect().await
45 })
46 }
47}
48
49pub struct ParAsyncFilter<I, F> {
51 iter: I,
52 concurrency: usize,
53 filter_fn: F,
54}
55
56impl<I, F> ParAsyncFilter<I, F> {
57 pub(super) fn new(iter: I, concurrency: usize, filter_fn: F) -> Self {
58 Self {
59 iter,
60 concurrency,
61 filter_fn,
62 }
63 }
64}
65
66impl<I, F, Fut> AsyncIterator for ParAsyncFilter<I, F>
67where
68 I: AsyncIterator,
69 F: Fn(&I::Item) -> Fut + Send + Sync,
70 Fut: Future<Output = bool> + Send,
71{
72 type Item = I::Item;
73
74 fn into_vec(self) -> Vec<Self::Item> {
75 let concurrency = self.concurrency.max(1);
76 let filter_fn = self.filter_fn;
77 let items = self.iter.into_vec();
78 futures::executor::block_on(async move {
79 let futures = stream::iter(items).map(|item| {
80 let filter_fn = &filter_fn;
81 async move {
82 let keep = filter_fn(&item).await;
83 (item, keep)
84 }
85 });
86 retained_buffered(futures, concurrency)
87 .filter_map(|(item, keep)| async move { keep.then_some(item) })
88 .collect()
89 .await
90 })
91 }
92}
93
94pub(super) async fn for_each<I, F, Fut>(iter: I, concurrency: usize, func: F)
101where
102 I: AsyncIterator,
103 F: Fn(I::Item) -> Fut + Send + Sync,
104 Fut: Future<Output = ()> + Send,
105{
106 let items = iter.into_vec();
107 retained_unordered(stream::iter(items).map(func), concurrency)
108 .for_each(|()| async {})
109 .await;
110}