Skip to main content

ParallelIterator

Trait ParallelIterator 

Source
pub trait ParallelIterator: Sized + Send {
    type Item: Send;

Show 76 methods // Required methods fn drive<C, R>(self, consumer: C) -> R where C: Consumer<Self::Item, Result = R> + Send + Sync, R: Send; fn seq_items(self) -> Vec<Self::Item>; // Provided methods fn seq_iter(self) -> impl Iterator<Item = Self::Item> { ... } fn seq_try_fold<T, B, F>(self, init: T, fold_fn: F) -> ControlFlow<B, T> where F: FnMut(T, Self::Item) -> ControlFlow<B, T> { ... } fn seq_fold<T, F>(self, init: T, fold_fn: F) -> T where F: FnMut(T, Self::Item) -> T { ... } fn seq_items_window( self, skip: usize, take: Option<usize>, ) -> Vec<Self::Item> { ... } fn seq_items_reversed(self) -> Vec<Self::Item> { ... } fn seq_items_reversed_prefix(self, count: usize) -> Vec<Self::Item> { ... } fn map<F, R>(self, map_fn: F) -> Map<Self, F> where F: Fn(Self::Item) -> R + Send + Sync + Clone, R: Send { ... } fn map_with<T, F, R>(self, init: T, map_fn: F) -> MapWith<Self, T, F> where T: Send + Clone, F: Fn(&mut T, Self::Item) -> R + Send + Sync + Clone, R: Send + Sync + 'static { ... } fn map_init<Init, T, F, R>( self, init: Init, map_fn: F, ) -> MapInit<Self, Init, F> where Init: Fn() -> T + Send + Sync + Clone, T: Send, F: Fn(&mut T, Self::Item) -> R + Send + Sync + Clone, R: Send + Sync + 'static { ... } fn update<F>(self, update_fn: F) -> Update<Self, F> where F: Fn(&mut Self::Item) + Send + Sync + Clone, Self::Item: Sync + 'static { ... } fn filter<F>(self, filter_fn: F) -> Filter<Self, F> where F: Fn(&Self::Item) -> bool + Send + Sync + Clone { ... } fn inspect<F>(self, inspect_fn: F) -> Inspect<Self, F> where F: Fn(&Self::Item) + Send + Sync + Clone, Self::Item: Sync { ... } fn panic_fuse(self) -> PanicFuse<Self> where Self::Item: Sync { ... } fn filter_map<F, R>(self, filter_map_fn: F) -> FilterMap<Self, F> where F: Fn(Self::Item) -> Option<R> + Send + Sync + Clone, R: Send + Sync + 'static { ... } fn while_some<T>(self) -> WhileSome<Self> where Self: ParallelIterator<Item = Option<T>>, T: Send + Sync + 'static { ... } fn flat_map<F, U>(self, flat_map_fn: F) -> FlatMap<Self, F> where F: Fn(Self::Item) -> U + Send + Sync + Clone, U: IntoIterator, <U as IntoIterator>::Item: Send + Sync + 'static { ... } fn flat_map_iter<F, U>(self, flat_map_fn: F) -> FlatMap<Self, F> where F: Fn(Self::Item) -> U + Send + Sync + Clone, U: IntoIterator, <U as IntoIterator>::Item: Send + Sync + 'static { ... } fn flatten(self) -> Flatten<Self> where Self::Item: IntoIterator, <Self::Item as IntoIterator>::Item: Send + Sync + 'static { ... } fn flatten_iter(self) -> Flatten<Self> where Self::Item: IntoIterator, <Self::Item as IntoIterator>::Item: Send + Sync + 'static { ... } fn enumerate(self) -> Enumerate<Self> where Self::Item: Sync + 'static { ... } fn zip<J>(self, other: J) -> Zip<Self, J> where J: ParallelIterator, Self::Item: Sync + 'static, <J as ParallelIterator>::Item: Sync + 'static { ... } fn zip_eq<J>(self, other: J) -> ZipEq<Self, J> where J: ParallelIterator, Self::Item: Sync + 'static, <J as ParallelIterator>::Item: Sync + 'static { ... } fn take(self, count: usize) -> Take<Self> where Self::Item: Sync + 'static { ... } fn take_any(self, count: usize) -> Take<Self> where Self::Item: Sync + 'static { ... } fn skip(self, count: usize) -> Skip<Self> where Self::Item: Sync + 'static { ... } fn skip_any(self, count: usize) -> Skip<Self> where Self::Item: Sync + 'static { ... } fn take_any_while<F>(self, predicate: F) -> TakeAnyWhile<Self, F> where F: Fn(&Self::Item) -> bool + Send + Sync + Clone, Self::Item: Sync + 'static { ... } fn skip_any_while<F>(self, predicate: F) -> SkipAnyWhile<Self, F> where F: Fn(&Self::Item) -> bool + Send + Sync + Clone, Self::Item: Sync + 'static { ... } fn chain<J>(self, other: J) -> Chain<Self, J> where J: ParallelIterator<Item = Self::Item>, Self::Item: Sync + 'static { ... } fn intersperse(self, separator: Self::Item) -> Intersperse<Self> where Self::Item: Clone + Sync + 'static { ... } fn rev(self) -> Rev<Self> where Self::Item: Sync + 'static { ... } fn chunks(self, chunk_size: usize) -> Chunks<Self> where Self::Item: Sync + 'static { ... } fn copied<'data, T>(self) -> Copied<Self> where Self: ParallelIterator<Item = &'data T>, T: Copy + Send + Sync + 'data + 'static { ... } fn cloned<'data, T>(self) -> Cloned<Self> where Self: ParallelIterator<Item = &'data T>, T: Clone + Send + Sync + 'data + 'static { ... } fn reduce<F>(self, reduce_fn: F) -> Option<Self::Item> where F: Fn(Self::Item, Self::Item) -> Self::Item + Send + Sync + Clone, Self::Item: Clone + Sync { ... } fn fold<T, F>(self, init: T, fold_fn: F) -> T where T: Send + Sync + Clone, F: Fn(T, Self::Item) -> T + Send + Sync + Clone, Self::Item: Sync { ... } fn collect<C>(self) -> C where C: ParallelExtend<Self::Item> + Default + Send { ... } fn collect_vec_list(self) -> LinkedList<Vec<Self::Item>> { ... } fn partition<C, F>(self, predicate: F) -> (C, C) where C: FromIterator<Self::Item> + Send, F: Fn(&Self::Item) -> bool + Send + Sync + Clone, Self::Item: Sync + 'static { ... } fn partition_map<A, B, P, L, R>(self, predicate: P) -> (A, B) where A: Default + Extend<L> + Send, B: Default + Extend<R> + Send, P: Fn(Self::Item) -> Either<L, R> + Send + Sync + Clone, L: Send, R: Send { ... } fn unzip<A, B, FromA, FromB>(self) -> (FromA, FromB) where Self: ParallelIterator<Item = (A, B)>, FromA: Default + Extend<A> + Send, FromB: Default + Extend<B> + Send, A: Send, B: Send { ... } fn sequential(self) -> SequentialAdapter<Self> { ... } fn count(self) -> usize where Self::Item: Sync { ... } fn find_first<F>(self, predicate: F) -> Option<Self::Item> where F: Fn(&Self::Item) -> bool + Send + Sync + Clone, Self::Item: Sync { ... } fn find_last<F>(self, predicate: F) -> Option<Self::Item> where F: Fn(&Self::Item) -> bool + Send + Sync + Clone { ... } fn position_first<F>(self, predicate: F) -> Option<usize> where F: Fn(Self::Item) -> bool + Send + Sync + Clone { ... } fn position_any<F>(self, predicate: F) -> Option<usize> where F: Fn(Self::Item) -> bool + Send + Sync + Clone { ... } fn position_last<F>(self, predicate: F) -> Option<usize> where F: Fn(Self::Item) -> bool + Send + Sync + Clone { ... } fn positions<F>(self, predicate: F) -> Positions<Self, F> where F: Fn(Self::Item) -> bool + Send + Sync + Clone { ... } fn find_map_first<F, R>(self, map_fn: F) -> Option<R> where F: Fn(Self::Item) -> Option<R> + Send + Sync + Clone, R: Send { ... } fn find_map_any<F, R>(self, map_fn: F) -> Option<R> where F: Fn(Self::Item) -> Option<R> + Send + Sync + Clone, R: Send { ... } fn find_map_last<F, R>(self, map_fn: F) -> Option<R> where F: Fn(Self::Item) -> Option<R> + Send + Sync + Clone, R: Send { ... } fn any<F>(self, predicate: F) -> bool where F: Fn(&Self::Item) -> bool + Send + Sync + Clone, Self::Item: Sync { ... } fn all<F>(self, predicate: F) -> bool where F: Fn(&Self::Item) -> bool + Send + Sync + Clone, Self::Item: Sync { ... } fn for_each<F>(self, op: F) where F: Fn(Self::Item) + Send + Sync + Clone { ... } fn for_each_with<T, F>(self, init: T, op: F) where T: Send + Clone, F: Fn(&mut T, Self::Item) + Send + Sync + Clone { ... } fn for_each_init<Init, T, F>(self, init: Init, op: F) where Init: Fn() -> T + Send + Sync + Clone, T: Send, F: Fn(&mut T, Self::Item) + Send + Sync + Clone { ... } fn try_for_each<F, E>(self, op: F) -> Result<(), E> where F: Fn(Self::Item) -> Result<(), E> + Send + Sync + Clone, E: Send { ... } fn try_for_each_with<T, F, E>(self, init: T, op: F) -> Result<(), E> where T: Send + Clone, F: Fn(&mut T, Self::Item) -> Result<(), E> + Send + Sync + Clone, E: Send { ... } fn try_for_each_init<Init, T, F, E>( self, init: Init, op: F, ) -> Result<(), E> where Init: Fn() -> T + Send + Sync + Clone, T: Send, F: Fn(&mut T, Self::Item) -> Result<(), E> + Send + Sync + Clone, E: Send { ... } fn reduce_with<F>(self, reduce_fn: F) -> Option<Self::Item> where F: Fn(Self::Item, Self::Item) -> Self::Item + Send + Sync + Clone, Self::Item: Sync + Clone { ... } fn try_reduce<Identity, F, T, E>( self, identity: Identity, reduce_fn: F, ) -> Result<T, E> where Self::Item: Into<Result<T, E>>, Identity: Fn() -> T + Send + Sync + Clone, F: Fn(T, T) -> Result<T, E> + Send + Sync + Clone, T: Send, E: Send { ... } fn try_reduce_with<F>(self, reduce_fn: F) -> Option<Self::Item> where Self::Item: TryStreamItem, F: Fn(<Self::Item as TryStreamItem>::Output, <Self::Item as TryStreamItem>::Output) -> Self::Item + Send + Sync + Clone { ... } fn sum<S>(self) -> S where S: Sum<Self::Item> + Send { ... } fn sum_reassociated<S>(self) -> S where S: Sum<Self::Item> + Sum + Send { ... } fn product<P>(self) -> P where P: Product<Self::Item> + Send { ... } fn product_reassociated<P>(self) -> P where P: Product<Self::Item> + Product + Send { ... } fn min(self) -> Option<Self::Item> where Self::Item: Ord { ... } fn max(self) -> Option<Self::Item> where Self::Item: Ord { ... } fn min_by<F>(self, compare: F) -> Option<Self::Item> where F: Fn(&Self::Item, &Self::Item) -> Ordering + Send + Sync + Clone { ... } fn max_by<F>(self, compare: F) -> Option<Self::Item> where F: Fn(&Self::Item, &Self::Item) -> Ordering + Send + Sync + Clone { ... } fn min_by_key<K, F>(self, key_fn: F) -> Option<Self::Item> where K: Ord, F: Fn(&Self::Item) -> K + Send + Sync + Clone { ... } fn max_by_key<K, F>(self, key_fn: F) -> Option<Self::Item> where K: Ord, F: Fn(&Self::Item) -> K + Send + Sync + Clone { ... } fn find_any<F>(self, predicate: F) -> Option<Self::Item> where F: Fn(&Self::Item) -> bool + Send + Sync + Clone, Self::Item: Sync { ... }
}
Expand description

Core parallel iterator trait for Moirai’s Rayon-style non-indexed subset.

Required Associated Types§

Source

type Item: Send

The type of items yielded by this parallel iterator.

Required Methods§

Source

fn drive<C, R>(self, consumer: C) -> R
where C: Consumer<Self::Item, Result = R> + Send + Sync, R: Send,

Drive the Consumer protocol over this iterator’s items.

§Concurrency contract

Large owned and borrowed vector sources split their consumer recursively and run one branch through Moirai’s nesting-safe SyncTask scope. Small shards remain inline so scheduler overhead does not dominate the work. Scope admission refusal runs the branch on the caller, preserving the every-item contract under shutdown or bounded-queue pressure. The resulting consumer combination preserves logical source order. The infallible iterator contract recovers an unclaimed branch on the caller if the scheduler cannot admit the scoped job; bounded admission refusal is handled by the scheduler’s caller-lane fallback before this method returns. A scheduler shutdown therefore degrades this drive to ordered caller-side execution rather than dropping work.

Source

fn seq_items(self) -> Vec<Self::Item>

Collect all items sequentially without routing through the consumer protocol.

Provided Methods§

Source

fn seq_iter(self) -> impl Iterator<Item = Self::Item>

Convert the logical item stream into a sequential iterator.

The default preserves compatibility for existing implementations by materializing through seq_items. Sources and adapters that can expose their logical stream directly override this method, allowing sequential terminals to retain one standard iterator invocation without allocating an intermediate vector.

§Examples
use moirai_iter::parallel::{IntoParallelIterator, ParallelIterator};

let items = vec![1_u32, 2, 3]
    .into_par_iter()
    .seq_iter()
    .collect::<Vec<_>>();
assert_eq!(items, vec![1, 2, 3]);
Source

fn seq_try_fold<T, B, F>(self, init: T, fold_fn: F) -> ControlFlow<B, T>
where F: FnMut(T, Self::Item) -> ControlFlow<B, T>,

Fold this iterator’s logical item stream left to right, stopping at the first ControlFlow::Break.

This is the folding counterpart to seq_iter and the base every folding Consumer runs on: a shard’s items reach the accumulator one at a time. The default delegates to seq_iter, whose compatibility implementation materializes through seq_items; sources and adapters on the terminal hot path override seq_iter to stream without an intermediate Vec.

The break value is the accumulator as it stood when the fold stopped, so a caller that needs the partial result on early exit reads it from the Break arm.

Source

fn seq_fold<T, F>(self, init: T, fold_fn: F) -> T
where F: FnMut(T, Self::Item) -> T,

Fold this iterator’s logical item stream left to right.

The non-short-circuiting form of seq_try_fold; it inherits that method’s streaming behaviour, so overriding either seq_iter or seq_try_fold is enough to make both allocation-free.

Source

fn seq_items_window(self, skip: usize, take: Option<usize>) -> Vec<Self::Item>

Collect a logical window from the sequential item stream.

Source

fn seq_items_reversed(self) -> Vec<Self::Item>

Collect items in reverse logical order.

Source

fn seq_items_reversed_prefix(self, count: usize) -> Vec<Self::Item>

Collect a prefix from the reversed logical item stream.

Source

fn map<F, R>(self, map_fn: F) -> Map<Self, F>
where F: Fn(Self::Item) -> R + Send + Sync + Clone, R: Send,

Map operation that transforms each element in parallel.

Source

fn map_with<T, F, R>(self, init: T, map_fn: F) -> MapWith<Self, T, F>
where T: Send + Clone, F: Fn(&mut T, Self::Item) -> R + Send + Sync + Clone, R: Send + Sync + 'static,

Map operation with cloned per-operation state.

Source

fn map_init<Init, T, F, R>( self, init: Init, map_fn: F, ) -> MapInit<Self, Init, F>
where Init: Fn() -> T + Send + Sync + Clone, T: Send, F: Fn(&mut T, Self::Item) -> R + Send + Sync + Clone, R: Send + Sync + 'static,

Map operation with lazily initialized state.

Source

fn update<F>(self, update_fn: F) -> Update<Self, F>
where F: Fn(&mut Self::Item) + Send + Sync + Clone, Self::Item: Sync + 'static,

Mutate each item by reference and yield the mutated item.

Source

fn filter<F>(self, filter_fn: F) -> Filter<Self, F>
where F: Fn(&Self::Item) -> bool + Send + Sync + Clone,

Filter operation that retains elements matching a predicate.

Source

fn inspect<F>(self, inspect_fn: F) -> Inspect<Self, F>
where F: Fn(&Self::Item) + Send + Sync + Clone, Self::Item: Sync,

Inspect each element by shared reference without changing the stream.

Source

fn panic_fuse(self) -> PanicFuse<Self>
where Self::Item: Sync,

Preserve value semantics while stopping sibling work after panic where applicable.

Source

fn filter_map<F, R>(self, filter_map_fn: F) -> FilterMap<Self, F>
where F: Fn(Self::Item) -> Option<R> + Send + Sync + Clone, R: Send + Sync + 'static,

Map each element to an optional value and retain present values.

Source

fn while_some<T>(self) -> WhileSome<Self>
where Self: ParallelIterator<Item = Option<T>>, T: Send + Sync + 'static,

Unwrap a prefix of present values from an optional stream.

Source

fn flat_map<F, U>(self, flat_map_fn: F) -> FlatMap<Self, F>
where F: Fn(Self::Item) -> U + Send + Sync + Clone, U: IntoIterator, <U as IntoIterator>::Item: Send + Sync + 'static,

Map each element to an iterator and flatten the resulting sequence.

Source

fn flat_map_iter<F, U>(self, flat_map_fn: F) -> FlatMap<Self, F>
where F: Fn(Self::Item) -> U + Send + Sync + Clone, U: IntoIterator, <U as IntoIterator>::Item: Send + Sync + 'static,

Map each element to a serial iterator and flatten the resulting sequence.

Source

fn flatten(self) -> Flatten<Self>
where Self::Item: IntoIterator, <Self::Item as IntoIterator>::Item: Send + Sync + 'static,

Flatten nested item streams with standard left-to-right semantics.

Source

fn flatten_iter(self) -> Flatten<Self>
where Self::Item: IntoIterator, <Self::Item as IntoIterator>::Item: Send + Sync + 'static,

Flatten nested serial iterators with standard left-to-right semantics.

Source

fn enumerate(self) -> Enumerate<Self>
where Self::Item: Sync + 'static,

Pair each element with its zero-based position in the logical sequence.

Source

fn zip<J>(self, other: J) -> Zip<Self, J>
where J: ParallelIterator, Self::Item: Sync + 'static, <J as ParallelIterator>::Item: Sync + 'static,

Pair elements with another parallel iterator, stopping at the shorter input.

Source

fn zip_eq<J>(self, other: J) -> ZipEq<Self, J>
where J: ParallelIterator, Self::Item: Sync + 'static, <J as ParallelIterator>::Item: Sync + 'static,

Pair elements with another parallel iterator and require equal lengths.

Source

fn take(self, count: usize) -> Take<Self>
where Self::Item: Sync + 'static,

Retain at most count elements from the logical sequence prefix.

Source

fn take_any(self, count: usize) -> Take<Self>
where Self::Item: Sync + 'static,

Retain at most count items from this non-indexed deterministic stream.

Source

fn skip(self, count: usize) -> Skip<Self>
where Self::Item: Sync + 'static,

Discard count elements from the logical sequence prefix.

Source

fn skip_any(self, count: usize) -> Skip<Self>
where Self::Item: Sync + 'static,

Discard count items from this non-indexed deterministic stream.

Source

fn take_any_while<F>(self, predicate: F) -> TakeAnyWhile<Self, F>
where F: Fn(&Self::Item) -> bool + Send + Sync + Clone, Self::Item: Sync + 'static,

Retain this deterministic stream prefix while predicate returns true.

Source

fn skip_any_while<F>(self, predicate: F) -> SkipAnyWhile<Self, F>
where F: Fn(&Self::Item) -> bool + Send + Sync + Clone, Self::Item: Sync + 'static,

Discard this deterministic stream prefix while predicate returns true.

Source

fn chain<J>(self, other: J) -> Chain<Self, J>
where J: ParallelIterator<Item = Self::Item>, Self::Item: Sync + 'static,

Concatenate this iterator with another iterator of the same item type.

Source

fn intersperse(self, separator: Self::Item) -> Intersperse<Self>
where Self::Item: Clone + Sync + 'static,

Insert a cloned separator between adjacent logical items.

Source

fn rev(self) -> Rev<Self>
where Self::Item: Sync + 'static,

Reverse the logical sequence order.

Source

fn chunks(self, chunk_size: usize) -> Chunks<Self>
where Self::Item: Sync + 'static,

Group the logical item stream into non-empty chunks.

Source

fn copied<'data, T>(self) -> Copied<Self>
where Self: ParallelIterator<Item = &'data T>, T: Copy + Send + Sync + 'data + 'static,

Copy referenced items out of a borrowed parallel stream.

Source

fn cloned<'data, T>(self) -> Cloned<Self>
where Self: ParallelIterator<Item = &'data T>, T: Clone + Send + Sync + 'data + 'static,

Clone referenced items out of a borrowed parallel stream.

Source

fn reduce<F>(self, reduce_fn: F) -> Option<Self::Item>
where F: Fn(Self::Item, Self::Item) -> Self::Item + Send + Sync + Clone, Self::Item: Clone + Sync,

Reduce operation that combines all elements.

Source

fn fold<T, F>(self, init: T, fold_fn: F) -> T
where T: Send + Sync + Clone, F: Fn(T, Self::Item) -> T + Send + Sync + Clone, Self::Item: Sync,

Fold operation with an initial value.

Source

fn collect<C>(self) -> C
where C: ParallelExtend<Self::Item> + Default + Send,

Collect into a collection.

Source

fn collect_vec_list(self) -> LinkedList<Vec<Self::Item>>

Collect into a list of owned vector segments.

This bounded terminal mirrors Rayon’s public collect_vec_list return shape while preserving Moirai’s logical item stream as one moved segment. Segment count is not part of the semantic contract; flattening the returned list yields the same logical item sequence as collect.

Source

fn partition<C, F>(self, predicate: F) -> (C, C)
where C: FromIterator<Self::Item> + Send, F: Fn(&Self::Item) -> bool + Send + Sync + Clone, Self::Item: Sync + 'static,

Partition items into two collections while preserving relative order.

Source

fn partition_map<A, B, P, L, R>(self, predicate: P) -> (A, B)
where A: Default + Extend<L> + Send, B: Default + Extend<R> + Send, P: Fn(Self::Item) -> Either<L, R> + Send + Sync + Clone, L: Send, R: Send,

Split mapped Either values into two collections while preserving side-local order.

Source

fn unzip<A, B, FromA, FromB>(self) -> (FromA, FromB)
where Self: ParallelIterator<Item = (A, B)>, FromA: Default + Extend<A> + Send, FromB: Default + Extend<B> + Send, A: Send, B: Send,

Split a stream of pairs into two collections while preserving order.

Source

fn sequential(self) -> SequentialAdapter<Self>

Convert to a sequential iterator.

Source

fn count(self) -> usize
where Self::Item: Sync,

Count the number of elements.

Source

fn find_first<F>(self, predicate: F) -> Option<Self::Item>
where F: Fn(&Self::Item) -> bool + Send + Sync + Clone, Self::Item: Sync,

Find the first element matching a predicate.

Every shard runs: a shard that has not started may hold an earlier match than one already found, so this terminal cannot abandon shards the way find_any does. Each shard still stops at its own first match.

Source

fn find_last<F>(self, predicate: F) -> Option<Self::Item>
where F: Fn(&Self::Item) -> bool + Send + Sync + Clone,

Find the last element matching a predicate in the logical stream.

Source

fn position_first<F>(self, predicate: F) -> Option<usize>
where F: Fn(Self::Item) -> bool + Send + Sync + Clone,

Find the first logical index matching a predicate.

Sequential by contract: a logical index is a property of the whole stream, and the non-indexed consumer protocol cannot hand a shard its own base index. Consumer::split_at carries the source split point, which a length-changing adapter such as filter invalidates before it reaches the shard. The stream is folded rather than collected, so no intermediate vector is built.

Source

fn position_any<F>(self, predicate: F) -> Option<usize>
where F: Fn(Self::Item) -> bool + Send + Sync + Clone,

Find any logical index matching a predicate.

Source

fn position_last<F>(self, predicate: F) -> Option<usize>
where F: Fn(Self::Item) -> bool + Send + Sync + Clone,

Find the last logical index matching a predicate.

Sequential for the reason given on position_first, and folded rather than collected.

Source

fn positions<F>(self, predicate: F) -> Positions<Self, F>
where F: Fn(Self::Item) -> bool + Send + Sync + Clone,

Return all logical indices whose items match a predicate.

Source

fn find_map_first<F, R>(self, map_fn: F) -> Option<R>
where F: Fn(Self::Item) -> Option<R> + Send + Sync + Clone, R: Send,

Find and map the first matching element in the logical stream.

Source

fn find_map_any<F, R>(self, map_fn: F) -> Option<R>
where F: Fn(Self::Item) -> Option<R> + Send + Sync + Clone, R: Send,

Find and map any matching element in the logical stream.

Shards that have not started are abandoned once any shard produces a mapped value, so the result is a mapped match rather than necessarily the logically first one. Use find_map_first when order matters.

Source

fn find_map_last<F, R>(self, map_fn: F) -> Option<R>
where F: Fn(Self::Item) -> Option<R> + Send + Sync + Clone, R: Send,

Find and map the last matching element in the logical stream.

Source

fn any<F>(self, predicate: F) -> bool
where F: Fn(&Self::Item) -> bool + Send + Sync + Clone, Self::Item: Sync,

Test if any element matches a predicate.

Source

fn all<F>(self, predicate: F) -> bool
where F: Fn(&Self::Item) -> bool + Send + Sync + Clone, Self::Item: Sync,

Test if all elements match a predicate.

Source

fn for_each<F>(self, op: F)
where F: Fn(Self::Item) + Send + Sync + Clone,

Apply a function to each element.

Source

fn for_each_with<T, F>(self, init: T, op: F)
where T: Send + Clone, F: Fn(&mut T, Self::Item) + Send + Sync + Clone,

Apply a function to each element with cloned per-operation state.

Sequential by contract: one state value threads through the whole stream, so op observes every prior item’s effect. A parallel form would have to give each shard its own clone, which is a different contract. The stream is folded rather than collected.

Source

fn for_each_init<Init, T, F>(self, init: Init, op: F)
where Init: Fn() -> T + Send + Sync + Clone, T: Send, F: Fn(&mut T, Self::Item) + Send + Sync + Clone,

Apply a function to each element with lazily initialized state.

Sequential for the reason given on for_each_with.

Source

fn try_for_each<F, E>(self, op: F) -> Result<(), E>
where F: Fn(Self::Item) -> Result<(), E> + Send + Sync + Clone, E: Send,

Apply a fallible function to each element and stop on the first error.

The returned error is the first one in logical order. Each shard stops at its own first error, but no shard is abandoned: an earlier shard may still hold an earlier error than one already reported.

Source

fn try_for_each_with<T, F, E>(self, init: T, op: F) -> Result<(), E>
where T: Send + Clone, F: Fn(&mut T, Self::Item) -> Result<(), E> + Send + Sync + Clone, E: Send,

Apply a fallible function to each element with cloned per-operation state.

Sequential for the reason given on for_each_with.

Source

fn try_for_each_init<Init, T, F, E>(self, init: Init, op: F) -> Result<(), E>
where Init: Fn() -> T + Send + Sync + Clone, T: Send, F: Fn(&mut T, Self::Item) -> Result<(), E> + Send + Sync + Clone, E: Send,

Apply a fallible function to each element with lazily initialized state.

Sequential for the reason given on for_each_with.

Source

fn reduce_with<F>(self, reduce_fn: F) -> Option<Self::Item>
where F: Fn(Self::Item, Self::Item) -> Self::Item + Send + Sync + Clone, Self::Item: Sync + Clone,

Reduce with an associative operation.

Source

fn try_reduce<Identity, F, T, E>( self, identity: Identity, reduce_fn: F, ) -> Result<T, E>
where Self::Item: Into<Result<T, E>>, Identity: Fn() -> T + Send + Sync + Clone, F: Fn(T, T) -> Result<T, E> + Send + Sync + Clone, T: Send, E: Send,

Reduce a fallible item stream with an identity and associative operation.

Source

fn try_reduce_with<F>(self, reduce_fn: F) -> Option<Self::Item>
where Self::Item: TryStreamItem, F: Fn(<Self::Item as TryStreamItem>::Output, <Self::Item as TryStreamItem>::Output) -> Self::Item + Send + Sync + Clone,

Reduce a fallible item stream without an identity value.

Source

fn sum<S>(self) -> S
where S: Sum<Self::Item> + Send,

Sum the complete logical stream through one standard Iterator::sum invocation.

std::iter::Sum does not expose an operation for combining partial output values. This method therefore preserves every lawful Sum<Self::Item> implementation through the iterator returned by seq_iter. Compatible sources and adapters stream directly; other implementations retain the default materialized path. Use sum_reassociated only when the output’s partial values may be reassociated.

Source

fn sum_reassociated<S>(self) -> S
where S: Sum<Self::Item> + Sum + Send,

Sum independently produced item fragments and merge their outputs.

This terminal invokes Sum<Self::Item> on empty and one-item streams, then invokes Sum<S> on pairs of partial outputs. That stronger contract enables parallel shard folding without materializing the full logical stream, but it is not equivalent to sum for an arbitrary Sum implementation.

§Ordering

Partial outputs are merged in logical shard order. The merge tree is a function of the input length alone, so arithmetic results are reproducible across runs and worker counts. Floating-point results need not be bit-identical to a strictly left-to-right sum.

Source

fn product<P>(self) -> P
where P: Product<Self::Item> + Send,

Multiply the complete logical stream through one standard Iterator::product invocation.

This preserves every lawful Product<Self::Item> implementation through seq_iter. Compatible sources and adapters stream directly; other implementations retain the default materialized path. Use product_reassociated only when partial output values may be reassociated.

Source

fn product_reassociated<P>(self) -> P
where P: Product<Self::Item> + Product + Send,

Multiply independently produced item fragments and merge their outputs.

This terminal invokes Product<Self::Item> on empty and one-item streams, then invokes Product<P> on pairs of partial outputs. See sum_reassociated for the deterministic merge ordering and semantic distinction from the standard terminal.

Source

fn min(self) -> Option<Self::Item>
where Self::Item: Ord,

Return the minimum item in the logical stream.

Source

fn max(self) -> Option<Self::Item>
where Self::Item: Ord,

Return the maximum item in the logical stream.

Source

fn min_by<F>(self, compare: F) -> Option<Self::Item>
where F: Fn(&Self::Item, &Self::Item) -> Ordering + Send + Sync + Clone,

Return the minimum item according to a comparator.

Ties resolve to the earliest item in logical order, matching Iterator::min_by. Shards keep their own earliest minimum and merges keep the earlier shard’s on equality, so the tie-break is the same at every level of the merge tree.

Source

fn max_by<F>(self, compare: F) -> Option<Self::Item>
where F: Fn(&Self::Item, &Self::Item) -> Ordering + Send + Sync + Clone,

Return the maximum item according to a comparator.

Ties resolve to the latest item in logical order, matching Iterator::max_by.

Source

fn min_by_key<K, F>(self, key_fn: F) -> Option<Self::Item>
where K: Ord, F: Fn(&Self::Item) -> K + Send + Sync + Clone,

Return the minimum item according to an ordered key.

Expressed through min_by, so tie-breaking matches Iterator::min_by_key. key_fn runs twice per comparison rather than being cached alongside the item, which keeps the key out of the value that crosses shard boundaries and so avoids a K: Send requirement.

Source

fn max_by_key<K, F>(self, key_fn: F) -> Option<Self::Item>
where K: Ord, F: Fn(&Self::Item) -> K + Send + Sync + Clone,

Return the maximum item according to an ordered key.

Expressed through max_by; see min_by_key for the key-evaluation note.

Source

fn find_any<F>(self, predicate: F) -> Option<Self::Item>
where F: Fn(&Self::Item) -> bool + Send + Sync + Clone, Self::Item: Sync,

Find any element matching a predicate.

Shards that have not started are abandoned once any shard finds a match, so the returned item is a match rather than necessarily the logically first one. Use find_first when order matters.

Dyn Compatibility§

This trait is not dyn compatible.

In older versions of Rust, dyn compatibility was called "object safety".

Implementors§

Source§

impl<'a, T> ParallelIterator for RefVecParIter<'a, T>
where T: Send + Sync,

Source§

impl<'data, I, T> ParallelIterator for Cloned<I>
where I: ParallelIterator<Item = &'data T>, T: Clone + Send + Sync + 'data + 'static,

Source§

type Item = T

Source§

impl<'data, I, T> ParallelIterator for Copied<I>
where I: ParallelIterator<Item = &'data T>, T: Copy + Send + Sync + 'data + 'static,

Source§

type Item = T

Source§

impl<'data, T> ParallelIterator for VecRefParIter<'data, T>
where T: Send + Sync + 'data,

Source§

impl<I, F, R> ParallelIterator for FilterMap<I, F>
where I: ParallelIterator, F: Fn(<I as ParallelIterator>::Item) -> Option<R> + Send + Sync + Clone, R: Send + Sync + 'static,

Source§

type Item = R

Source§

impl<I, F, R> ParallelIterator for Map<I, F>
where I: ParallelIterator, F: Fn(<I as ParallelIterator>::Item) -> R + Send + Sync + Clone, R: Send,

Source§

type Item = R

Source§

impl<I, F, U> ParallelIterator for FlatMap<I, F>
where I: ParallelIterator, F: Fn(<I as ParallelIterator>::Item) -> U + Send + Sync + Clone, U: IntoIterator, <U as IntoIterator>::Item: Send + Sync + 'static,

Source§

impl<I, F> ParallelIterator for Filter<I, F>
where I: ParallelIterator, F: Fn(&<I as ParallelIterator>::Item) -> bool + Send + Sync + Clone,

Source§

impl<I, F> ParallelIterator for Inspect<I, F>
where I: ParallelIterator, F: Fn(&<I as ParallelIterator>::Item) + Send + Sync + Clone, <I as ParallelIterator>::Item: Sync + 'static,

Source§

impl<I, F> ParallelIterator for Positions<I, F>
where I: ParallelIterator, F: Fn(<I as ParallelIterator>::Item) -> bool + Send + Sync + Clone,

Source§

impl<I, F> ParallelIterator for SkipAnyWhile<I, F>
where I: ParallelIterator, F: Fn(&<I as ParallelIterator>::Item) -> bool + Send + Sync + Clone, <I as ParallelIterator>::Item: Sync + 'static,

Source§

impl<I, F> ParallelIterator for TakeAnyWhile<I, F>
where I: ParallelIterator, F: Fn(&<I as ParallelIterator>::Item) -> bool + Send + Sync + Clone, <I as ParallelIterator>::Item: Sync + 'static,

Source§

impl<I, F> ParallelIterator for Update<I, F>
where I: ParallelIterator, F: Fn(&mut <I as ParallelIterator>::Item) + Send + Sync + Clone, <I as ParallelIterator>::Item: Sync + 'static,

Source§

impl<I, Init, T, F, R> ParallelIterator for MapInit<I, Init, F>
where I: ParallelIterator, Init: Fn() -> T + Send + Sync + Clone, T: Send, F: Fn(&mut T, <I as ParallelIterator>::Item) -> R + Send + Sync + Clone, R: Send + Sync + 'static,

Source§

type Item = R

Source§

impl<I, J> ParallelIterator for Chain<I, J>
where I: ParallelIterator, J: ParallelIterator<Item = <I as ParallelIterator>::Item>, <I as ParallelIterator>::Item: Sync + 'static,

Source§

impl<I, J> ParallelIterator for Interleave<I, J>
where I: ParallelIterator, J: ParallelIterator<Item = <I as ParallelIterator>::Item>, <I as ParallelIterator>::Item: Sync + 'static,

Source§

impl<I, J> ParallelIterator for InterleaveShortest<I, J>
where I: ParallelIterator, J: ParallelIterator<Item = <I as ParallelIterator>::Item>, <I as ParallelIterator>::Item: Sync + 'static,

Source§

impl<I, J> ParallelIterator for Zip<I, J>
where I: ParallelIterator, J: ParallelIterator, <I as ParallelIterator>::Item: Sync + 'static, <J as ParallelIterator>::Item: Sync + 'static,

Source§

impl<I, J> ParallelIterator for ZipEq<I, J>
where I: ParallelIterator, J: ParallelIterator, <I as ParallelIterator>::Item: Sync + 'static, <J as ParallelIterator>::Item: Sync + 'static,

Source§

impl<I, MapFn, Predicate, Mapped> ParallelIterator for MapPositions<I, MapFn, Predicate>
where I: ParallelIterator, MapFn: Fn(<I as ParallelIterator>::Item) -> Mapped + Send + Sync + Clone, Predicate: Fn(Mapped) -> bool + Send + Sync + Clone, Mapped: Send,

Source§

impl<I, T, F, R> ParallelIterator for MapWith<I, T, F>
where I: ParallelIterator, T: Send + Clone, F: Fn(&mut T, <I as ParallelIterator>::Item) -> R + Send + Sync + Clone, R: Send + Sync + 'static,

Source§

type Item = R

Source§

impl<I, T> ParallelIterator for WhileSome<I>
where I: ParallelIterator<Item = Option<T>>, T: Send + Sync + 'static,

Source§

type Item = T

Source§

impl<I> ParallelIterator for Chunks<I>
where I: ParallelIterator, <I as ParallelIterator>::Item: Sync + 'static,

Source§

impl<I> ParallelIterator for Enumerate<I>
where I: ParallelIterator, <I as ParallelIterator>::Item: Sync + 'static,

Source§

impl<I> ParallelIterator for ExponentialBlocks<I>
where I: ParallelIterator, <I as ParallelIterator>::Item: Sync + 'static,

Source§

impl<I> ParallelIterator for Flatten<I>

Source§

impl<I> ParallelIterator for Intersperse<I>
where I: ParallelIterator, <I as ParallelIterator>::Item: Clone + Sync + 'static,

Source§

impl<I> ParallelIterator for PanicFuse<I>

Source§

impl<I> ParallelIterator for Rev<I>
where I: ParallelIterator, <I as ParallelIterator>::Item: Sync + 'static,

Source§

impl<I> ParallelIterator for SequentialIterAdapter<I>
where I: Iterator + Send, <I as Iterator>::Item: Send + Sync + 'static,

Source§

type Item = <I as Iterator>::Item

Source§

impl<I> ParallelIterator for Skip<I>
where I: ParallelIterator, <I as ParallelIterator>::Item: Sync + 'static,

Source§

impl<I> ParallelIterator for StepBy<I>
where I: ParallelIterator, <I as ParallelIterator>::Item: Sync + 'static,

Source§

impl<I> ParallelIterator for Take<I>
where I: ParallelIterator, <I as ParallelIterator>::Item: Sync + 'static,

Source§

impl<I> ParallelIterator for UniformBlocks<I>
where I: ParallelIterator, <I as ParallelIterator>::Item: Sync + 'static,

Source§

impl<T> ParallelIterator for RangeParIter<T>
where T: Send + Sync + Clone + 'static + PartialOrd + Add<Output = T> + From<u8>,

Source§

type Item = T

Source§

impl<T> ParallelIterator for VecParIter<T>
where T: Send + Sync + 'static,

Source§

type Item = T