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