bevy_ecs/query/par_iter.rs
1use crate::{
2 batching::BatchingStrategy,
3 change_detection::Tick,
4 entity::{EntityEquivalent, UniqueEntityEquivalentVec},
5 query::{ArchetypeFilter, ContiguousQueryData, QueryContiguousIter, QueryEntityError},
6 world::unsafe_world_cell::UnsafeWorldCell,
7};
8
9use super::{IterQueryData, QueryFilter, QueryItem, QueryState, ReadOnlyQueryData};
10
11use alloc::vec::Vec;
12
13/// A parallel iterator over query results of a [`Query`](crate::system::Query).
14///
15/// This struct is created by the [`Query::par_iter`](crate::system::Query::par_iter) and
16/// [`Query::par_iter_mut`](crate::system::Query::par_iter_mut) methods.
17pub struct QueryParIter<'w, 's, D: IterQueryData, F: QueryFilter> {
18 pub(crate) world: UnsafeWorldCell<'w>,
19 pub(crate) state: &'s QueryState<D, F>,
20 pub(crate) last_run: Tick,
21 pub(crate) this_run: Tick,
22 pub(crate) batching_strategy: BatchingStrategy,
23}
24
25impl<'w, 's, D: IterQueryData, F: QueryFilter> QueryParIter<'w, 's, D, F> {
26 /// Changes the batching strategy used when iterating.
27 ///
28 /// For more information on how this affects the resultant iteration, see
29 /// [`BatchingStrategy`].
30 pub fn batching_strategy(mut self, strategy: BatchingStrategy) -> Self {
31 self.batching_strategy = strategy;
32 self
33 }
34
35 /// Runs `func` on each query result in parallel.
36 ///
37 /// # Panics
38 /// If the [`ComputeTaskPool`] is not initialized. If using this from a query that is being
39 /// initialized and run from the ECS scheduler, this should never panic.
40 ///
41 /// [`ComputeTaskPool`]: bevy_tasks::ComputeTaskPool
42 #[inline]
43 pub fn for_each<FN: Fn(QueryItem<'w, 's, D>) + Send + Sync + Clone>(self, func: FN) {
44 self.for_each_init(|| {}, |_, item| func(item));
45 }
46
47 /// Runs `func` on each query result in parallel on a value returned by `init`.
48 ///
49 /// `init` may be called multiple times per thread, and the values returned may be discarded between tasks on any given thread.
50 /// Callers should avoid using this function as if it were a parallel version
51 /// of [`Iterator::fold`].
52 ///
53 /// # Example
54 ///
55 /// ```
56 /// use bevy_utils::Parallel;
57 /// use crate::{bevy_ecs::prelude::Component, bevy_ecs::system::Query};
58 /// #[derive(Component)]
59 /// struct T;
60 /// fn system(query: Query<&T>){
61 /// let mut queue: Parallel<usize> = Parallel::default();
62 /// // queue.borrow_local_mut() will get or create a thread_local queue for each task/thread;
63 /// query.par_iter().for_each_init(|| queue.borrow_local_mut(),|local_queue, item| {
64 /// **local_queue += 1;
65 /// });
66 ///
67 /// // collect value from every thread
68 /// let entity_count: usize = queue.iter_mut().map(|v| *v).sum();
69 /// }
70 /// ```
71 ///
72 /// # Panics
73 /// If the [`ComputeTaskPool`] is not initialized. If using this from a query that is being
74 /// initialized and run from the ECS scheduler, this should never panic.
75 ///
76 /// [`ComputeTaskPool`]: bevy_tasks::ComputeTaskPool
77 #[inline]
78 pub fn for_each_init<FN, INIT, T>(self, init: INIT, func: FN)
79 where
80 FN: Fn(&mut T, QueryItem<'w, 's, D>) + Send + Sync + Clone,
81 INIT: Fn() -> T + Sync + Send + Clone,
82 {
83 let func = |mut init, item| {
84 func(&mut init, item);
85 init
86 };
87 #[cfg(any(target_arch = "wasm32", not(feature = "multi_threaded")))]
88 {
89 let init = init();
90 // SAFETY:
91 // This method can only be called once per instance of QueryParIter,
92 // which ensures that mutable queries cannot be executed multiple times at once.
93 // Mutable instances of QueryParIter can only be created via an exclusive borrow of a
94 // Query or a World, which ensures that multiple aliasing QueryParIters cannot exist
95 // at the same time.
96 unsafe {
97 self.state
98 .query_unchecked_manual_with_ticks(self.world, self.last_run, self.this_run)
99 .into_iter()
100 .fold(init, func);
101 }
102 }
103 #[cfg(all(not(target_arch = "wasm32"), feature = "multi_threaded"))]
104 {
105 let thread_count = bevy_tasks::ComputeTaskPool::get().thread_num();
106 if thread_count <= 1 {
107 let init = init();
108 // SAFETY: See the safety comment above.
109 unsafe {
110 self.state
111 .query_unchecked_manual_with_ticks(self.world, self.last_run, self.this_run)
112 .into_iter()
113 .fold(init, func);
114 }
115 } else {
116 // Need a batch size of at least 1.
117 let batch_size = self.get_batch_size(thread_count).max(1);
118 // SAFETY: See the safety comment above.
119 unsafe {
120 self.state.par_fold_init_unchecked_manual(
121 init,
122 self.world,
123 batch_size,
124 func,
125 self.last_run,
126 self.this_run,
127 );
128 }
129 }
130 }
131 }
132
133 #[cfg(all(not(target_arch = "wasm32"), feature = "multi_threaded"))]
134 fn get_batch_size(&self, thread_count: usize) -> u32 {
135 let max_items = || {
136 let id_iter = self.state.matched_storage_ids.iter();
137 if self.state.is_dense {
138 // SAFETY: We only access table metadata.
139 let tables = unsafe { &self.world.world_metadata().storages().tables };
140 id_iter
141 // SAFETY: The if check ensures that matched_storage_ids stores TableIds
142 .map(|id| unsafe { tables[id.table_id].entity_count() })
143 .max()
144 } else {
145 let archetypes = &self.world.archetypes();
146 id_iter
147 // SAFETY: The if check ensures that matched_storage_ids stores ArchetypeIds
148 .map(|id| unsafe { archetypes[id.archetype_id].len() })
149 .max()
150 }
151 .map(|v| v as usize)
152 .unwrap_or(0)
153 };
154 self.batching_strategy
155 .calc_batch_size(max_items, thread_count) as u32
156 }
157}
158
159/// A parallel iterator over contiguous query results on a
160/// [`Query`](crate::system::Query).
161///
162/// The
163/// [`Query::contiguous_par_iter`](crate::system::Query::contiguous_par_iter)
164/// and
165/// [`Query::contiguous_par_iter_mut`](crate::system::Query::contiguous_par_iter_mut)
166/// methods create instances of this structure.
167pub struct QueryContiguousParIter<'w, 's, D, F>
168where
169 D: ContiguousQueryData,
170 F: ArchetypeFilter,
171{
172 /// A reference to the world that contains the components that this query
173 /// iterates over.
174 pub(crate) world: UnsafeWorldCell<'w>,
175 /// Scoped access to the world state.
176 pub(crate) state: &'s QueryState<D, F>,
177 /// The tick that corresponds to the previous time this query ran.
178 pub(crate) last_run: Tick,
179 /// The tick that corresponds to the current run of the query.
180 pub(crate) this_run: Tick,
181 /// How matched rows are to be divided among worker threads.
182 pub(crate) batching_strategy: BatchingStrategy,
183}
184
185impl<'w, 's, D, F> QueryContiguousParIter<'w, 's, D, F>
186where
187 D: ContiguousQueryData,
188 F: ArchetypeFilter,
189{
190 /// Returns `None` if `query_state` is not dense, and hence not contiguously iterable.
191 pub(crate) fn new(
192 world: UnsafeWorldCell<'w>,
193 state: &'s QueryState<D, F>,
194 last_run: Tick,
195 this_run: Tick,
196 ) -> Option<Self> {
197 state.is_dense.then(|| Self {
198 world,
199 state,
200 last_run,
201 this_run,
202 batching_strategy: BatchingStrategy::new(),
203 })
204 }
205
206 /// Changes the batching strategy used when iterating.
207 ///
208 /// For more information on how this affects the resultant iteration, see
209 /// [`BatchingStrategy`].
210 pub fn batching_strategy(mut self, strategy: BatchingStrategy) -> Self {
211 self.batching_strategy = strategy;
212 self
213 }
214
215 /// Runs `func` on each contiguous chunk of query results in parallel.
216 ///
217 /// # Panics
218 /// If the [`ComputeTaskPool`] is not initialized. If using this from a
219 /// query that is being initialized and run from the ECS scheduler, this
220 /// should never panic.
221 ///
222 /// [`ComputeTaskPool`]: bevy_tasks::ComputeTaskPool
223 #[inline]
224 pub fn for_each(self, func: impl Fn(D::Contiguous<'w, 's>) + Send + Sync + Clone) {
225 self.for_each_init(|| {}, |_, item| func(item));
226 }
227
228 /// Runs `func` on each query result in parallel on a value returned by `init`.
229 ///
230 /// `init` may be called multiple times per thread, and the values returned may be discarded between tasks on any given thread.
231 /// Callers should avoid using this function as if it were a parallel version
232 /// of [`Iterator::fold`].
233 ///
234 /// # Example
235 ///
236 /// ```
237 /// use bevy_utils::Parallel;
238 /// use crate::{bevy_ecs::prelude::Component, bevy_ecs::system::Query};
239 /// #[derive(Component)]
240 /// struct T;
241 /// fn system(query: Query<&T>){
242 /// let mut queue: Parallel<usize> = Parallel::default();
243 /// // queue.borrow_local_mut() will get or create a thread_local queue for each task/thread.
244 /// // We unwrap the call to `contiguous_par_iter()` because we know the query in question is dense.
245 /// query.contiguous_par_iter().unwrap().for_each_init(|| queue.borrow_local_mut(),|local_queue, items| {
246 /// for _ in items {
247 /// **local_queue += 1;
248 /// }
249 /// });
250 ///
251 /// // collect value from every thread
252 /// let entity_count: usize = queue.iter_mut().map(|v| *v).sum();
253 /// }
254 /// ```
255 ///
256 /// # Panics
257 /// If the [`ComputeTaskPool`] is not initialized. If using this from a
258 /// query that is being initialized and run from the ECS scheduler, this
259 /// should never panic.
260 ///
261 /// [`ComputeTaskPool`]: bevy_tasks::ComputeTaskPool
262 pub fn for_each_init<T>(
263 self,
264 init: impl Fn() -> T + Sync + Send + Clone,
265 func: impl Fn(&mut T, D::Contiguous<'w, 's>) + Send + Sync + Clone,
266 ) {
267 let func = |mut init, item| {
268 func(&mut init, item);
269 init
270 };
271
272 #[cfg(any(target_arch = "wasm32", not(feature = "multi_threaded")))]
273 unsafe {
274 // SAFETY: This method can only be called once per instance of
275 // `QueryContiguousParIter`, which ensures that mutable queries
276 // cannot be executed multiple times at once. Mutable instances of
277 // `QueryContiguousParIter` can only be created via an exclusive
278 // borrow of a `Query` or a `World`, which ensures that multiple
279 // aliasing `QueryContiguousParIter`s cannot exist at the same time.
280 QueryContiguousIter::new(self.world, self.state, self.last_run, self.this_run)
281 .unwrap()
282 .fold(init(), func);
283 }
284
285 #[cfg(all(not(target_arch = "wasm32"), feature = "multi_threaded"))]
286 {
287 let thread_count = bevy_tasks::ComputeTaskPool::get().thread_num();
288 // SAFETY: This method can only be called once per instance of
289 // `QueryContiguousParIter`, which ensures that mutable queries
290 // cannot be executed multiple times at once. Mutable instances of
291 // `QueryContiguousParIter` can only be created via an exclusive
292 // borrow of a `Query` or a `World`, which ensures that multiple
293 // aliasing `QueryContiguousParIter`s cannot exist at the same time.
294 unsafe {
295 if thread_count <= 1 {
296 // Just run sequentially.
297 QueryContiguousIter::new(self.world, self.state, self.last_run, self.this_run)
298 .unwrap()
299 .fold(init(), func);
300 return;
301 }
302
303 // Dispatch to `contiguous_par_fold_init_unchecked_manual` for
304 // parallel iteration.
305 let batch_size = self.get_batch_size(thread_count).max(1);
306 self.state.contiguous_par_fold_init_unchecked_manual(
307 init,
308 self.world,
309 batch_size,
310 func,
311 self.last_run,
312 self.this_run,
313 );
314 }
315 }
316 }
317
318 /// Returns the size of each batch in rows, given a thread count and the
319 /// current batching strategy.
320 #[cfg(all(not(target_arch = "wasm32"), feature = "multi_threaded"))]
321 fn get_batch_size(&self, thread_count: usize) -> u32 {
322 let max_items = || {
323 let id_iter = self.state.matched_storage_ids.iter();
324 // SAFETY: We only access table metadata.
325 let tables = unsafe { &self.world.storages().tables };
326 id_iter
327 .map(|id| {
328 // SAFETY: Contiguous iteration can only process tables, so
329 // we must have a table here.
330 let table_id = unsafe { id.table_id };
331 tables[table_id].entity_count()
332 })
333 .max()
334 .map(|v| v as usize)
335 .unwrap_or(0)
336 };
337 self.batching_strategy
338 .calc_batch_size(max_items, thread_count) as u32
339 }
340}
341
342/// A parallel iterator over the unique query items generated from an [`Entity`] list.
343///
344/// This struct is created by the [`Query::par_iter_many`] method.
345///
346/// [`Entity`]: crate::entity::Entity
347/// [`Query::par_iter_many`]: crate::system::Query::par_iter_many
348pub struct QueryParManyIter<'w, 's, D: IterQueryData, F: QueryFilter, E: EntityEquivalent> {
349 pub(crate) world: UnsafeWorldCell<'w>,
350 pub(crate) state: &'s QueryState<D, F>,
351 pub(crate) entity_list: Vec<E>,
352 pub(crate) last_run: Tick,
353 pub(crate) this_run: Tick,
354 pub(crate) batching_strategy: BatchingStrategy,
355}
356
357impl<'w, 's, D: ReadOnlyQueryData, F: QueryFilter, E: EntityEquivalent + Sync>
358 QueryParManyIter<'w, 's, D, F, E>
359{
360 /// Changes the batching strategy used when iterating.
361 ///
362 /// For more information on how this affects the resultant iteration, see
363 /// [`BatchingStrategy`].
364 pub fn batching_strategy(mut self, strategy: BatchingStrategy) -> Self {
365 self.batching_strategy = strategy;
366 self
367 }
368
369 /// Runs `func` on each query result in parallel.
370 ///
371 /// # Panics
372 /// If the [`ComputeTaskPool`] is not initialized. If using this from a query that is being
373 /// initialized and run from the ECS scheduler, this should never panic.
374 ///
375 /// [`ComputeTaskPool`]: bevy_tasks::ComputeTaskPool
376 #[inline]
377 pub fn for_each<
378 FN: Fn(Result<QueryItem<'w, 's, D>, QueryEntityError>) + Send + Sync + Clone,
379 >(
380 self,
381 func: FN,
382 ) {
383 self.for_each_init(|| {}, |_, item| func(item));
384 }
385
386 /// Runs `func` on each query result in parallel on a value returned by `init`.
387 ///
388 /// `init` may be called multiple times per thread, and the values returned may be discarded between tasks on any given thread.
389 /// Callers should avoid using this function as if it were a parallel version
390 /// of [`Iterator::fold`].
391 ///
392 /// # Example
393 ///
394 /// ```
395 /// use bevy_utils::Parallel;
396 /// use crate::{bevy_ecs::prelude::{Component, Res, Resource, Entity}, bevy_ecs::system::Query};
397 /// # use core::slice;
398 /// use bevy_platform::prelude::Vec;
399 /// # fn some_expensive_operation(_item: &T) -> usize {
400 /// # 0
401 /// # }
402 ///
403 /// #[derive(Component)]
404 /// struct T;
405 ///
406 /// #[derive(Resource)]
407 /// struct V(Vec<Entity>);
408 ///
409 /// impl<'a> IntoIterator for &'a V {
410 /// // ...
411 /// # type Item = &'a Entity;
412 /// # type IntoIter = slice::Iter<'a, Entity>;
413 /// #
414 /// # fn into_iter(self) -> Self::IntoIter {
415 /// # self.0.iter()
416 /// # }
417 /// }
418 ///
419 /// fn system(query: Query<&T>, entities: Res<V>){
420 /// let mut queue: Parallel<usize> = Parallel::default();
421 /// // queue.borrow_local_mut() will get or create a thread_local queue for each task/thread;
422 /// query.par_iter_many(&entities).for_each_init(|| queue.borrow_local_mut(),|local_queue, item| {
423 /// **local_queue += some_expensive_operation(item.unwrap());
424 /// });
425 ///
426 /// // collect value from every thread
427 /// let final_value: usize = queue.iter_mut().map(|v| *v).sum();
428 /// }
429 /// # bevy_ecs::system::assert_is_system(system);
430 /// ```
431 ///
432 /// # Panics
433 /// If the [`ComputeTaskPool`] is not initialized. If using this from a query that is being
434 /// initialized and run from the ECS scheduler, this should never panic.
435 ///
436 /// [`ComputeTaskPool`]: bevy_tasks::ComputeTaskPool
437 #[inline]
438 pub fn for_each_init<FN, INIT, T>(self, init: INIT, func: FN)
439 where
440 FN: Fn(&mut T, Result<QueryItem<'w, 's, D>, QueryEntityError>) + Send + Sync + Clone,
441 INIT: Fn() -> T + Sync + Send + Clone,
442 {
443 let func = |mut init, item| {
444 func(&mut init, item);
445 init
446 };
447 #[cfg(any(target_arch = "wasm32", not(feature = "multi_threaded")))]
448 {
449 let init = init();
450 // SAFETY:
451 // This method can only be called once per instance of QueryParManyIter,
452 // which ensures that mutable queries cannot be executed multiple times at once.
453 // Mutable instances of QueryParManyUniqueIter can only be created via an exclusive borrow of a
454 // Query or a World, which ensures that multiple aliasing QueryParManyIters cannot exist
455 // at the same time.
456 unsafe {
457 self.state
458 .query_unchecked_manual_with_ticks(self.world, self.last_run, self.this_run)
459 .iter_many_inner(&self.entity_list)
460 .fold(init, func);
461 }
462 }
463 #[cfg(all(not(target_arch = "wasm32"), feature = "multi_threaded"))]
464 {
465 let thread_count = bevy_tasks::ComputeTaskPool::get().thread_num();
466 if thread_count <= 1 {
467 let init = init();
468 // SAFETY: See the safety comment above.
469 unsafe {
470 self.state
471 .query_unchecked_manual_with_ticks(self.world, self.last_run, self.this_run)
472 .iter_many_inner(&self.entity_list)
473 .fold(init, func);
474 }
475 } else {
476 // Need a batch size of at least 1.
477 let batch_size = self.get_batch_size(thread_count).max(1);
478 // SAFETY: See the safety comment above.
479 unsafe {
480 self.state.par_many_fold_init_unchecked_manual(
481 init,
482 self.world,
483 &self.entity_list,
484 batch_size,
485 func,
486 self.last_run,
487 self.this_run,
488 );
489 }
490 }
491 }
492 }
493
494 #[cfg(all(not(target_arch = "wasm32"), feature = "multi_threaded"))]
495 fn get_batch_size(&self, thread_count: usize) -> u32 {
496 self.batching_strategy
497 .calc_batch_size(|| self.entity_list.len(), thread_count) as u32
498 }
499}
500
501/// A parallel iterator over the unique query items generated from an [`EntitySet`].
502///
503/// This struct is created by the [`Query::par_iter_many_unique`] and [`Query::par_iter_many_unique_mut`] methods.
504///
505/// [`EntitySet`]: crate::entity::EntitySet
506/// [`Query::par_iter_many_unique`]: crate::system::Query::par_iter_many_unique
507/// [`Query::par_iter_many_unique_mut`]: crate::system::Query::par_iter_many_unique_mut
508pub struct QueryParManyUniqueIter<
509 'w,
510 's,
511 D: IterQueryData,
512 F: QueryFilter,
513 E: EntityEquivalent + Sync,
514> {
515 pub(crate) world: UnsafeWorldCell<'w>,
516 pub(crate) state: &'s QueryState<D, F>,
517 pub(crate) entity_list: UniqueEntityEquivalentVec<E>,
518 pub(crate) last_run: Tick,
519 pub(crate) this_run: Tick,
520 pub(crate) batching_strategy: BatchingStrategy,
521}
522
523impl<'w, 's, D: IterQueryData, F: QueryFilter, E: EntityEquivalent + Sync>
524 QueryParManyUniqueIter<'w, 's, D, F, E>
525{
526 /// Changes the batching strategy used when iterating.
527 ///
528 /// For more information on how this affects the resultant iteration, see
529 /// [`BatchingStrategy`].
530 pub fn batching_strategy(mut self, strategy: BatchingStrategy) -> Self {
531 self.batching_strategy = strategy;
532 self
533 }
534
535 /// Runs `func` on each query result in parallel.
536 ///
537 /// # Panics
538 /// If the [`ComputeTaskPool`] is not initialized. If using this from a query that is being
539 /// initialized and run from the ECS scheduler, this should never panic.
540 ///
541 /// [`ComputeTaskPool`]: bevy_tasks::ComputeTaskPool
542 #[inline]
543 pub fn for_each<
544 FN: Fn(Result<QueryItem<'w, 's, D>, QueryEntityError>) + Send + Sync + Clone,
545 >(
546 self,
547 func: FN,
548 ) {
549 self.for_each_init(|| {}, |_, item| func(item));
550 }
551
552 /// Runs `func` on each query result in parallel on a value returned by `init`.
553 ///
554 /// `init` may be called multiple times per thread, and the values returned may be discarded between tasks on any given thread.
555 /// Callers should avoid using this function as if it were a parallel version
556 /// of [`Iterator::fold`].
557 ///
558 /// # Example
559 ///
560 /// ```
561 /// use bevy_utils::Parallel;
562 /// use crate::{bevy_ecs::{prelude::{Component, Res, Resource, Entity}, entity::UniqueEntityVec, system::Query}};
563 /// # use core::slice;
564 /// # use crate::bevy_ecs::entity::UniqueEntityIter;
565 /// # fn some_expensive_operation(_item: &T) -> usize {
566 /// # 0
567 /// # }
568 ///
569 /// #[derive(Component)]
570 /// struct T;
571 ///
572 /// #[derive(Resource)]
573 /// struct V(UniqueEntityVec);
574 ///
575 /// impl<'a> IntoIterator for &'a V {
576 /// // ...
577 /// # type Item = &'a Entity;
578 /// # type IntoIter = UniqueEntityIter<slice::Iter<'a, Entity>>;
579 /// #
580 /// # fn into_iter(self) -> Self::IntoIter {
581 /// # self.0.iter()
582 /// # }
583 /// }
584 ///
585 /// fn system(query: Query<&T>, entities: Res<V>){
586 /// let mut queue: Parallel<usize> = Parallel::default();
587 /// // queue.borrow_local_mut() will get or create a thread_local queue for each task/thread;
588 /// query.par_iter_many_unique(&entities).for_each_init(|| queue.borrow_local_mut(),|local_queue, item| {
589 /// **local_queue += some_expensive_operation(item.unwrap());
590 /// });
591 ///
592 /// // collect value from every thread
593 /// let final_value: usize = queue.iter_mut().map(|v| *v).sum();
594 /// }
595 /// # bevy_ecs::system::assert_is_system(system);
596 /// ```
597 ///
598 /// # Panics
599 /// If the [`ComputeTaskPool`] is not initialized. If using this from a query that is being
600 /// initialized and run from the ECS scheduler, this should never panic.
601 ///
602 /// [`ComputeTaskPool`]: bevy_tasks::ComputeTaskPool
603 #[inline]
604 pub fn for_each_init<FN, INIT, T>(self, init: INIT, func: FN)
605 where
606 FN: Fn(&mut T, Result<QueryItem<'w, 's, D>, QueryEntityError>) + Send + Sync + Clone,
607 INIT: Fn() -> T + Sync + Send + Clone,
608 {
609 let func = |mut init, item| {
610 func(&mut init, item);
611 init
612 };
613 #[cfg(any(target_arch = "wasm32", not(feature = "multi_threaded")))]
614 {
615 let init = init();
616 // SAFETY:
617 // This method can only be called once per instance of QueryParManyUniqueIter,
618 // which ensures that mutable queries cannot be executed multiple times at once.
619 // Mutable instances of QueryParManyUniqueIter can only be created via an exclusive borrow of a
620 // Query or a World, which ensures that multiple aliasing QueryParManyUniqueIters cannot exist
621 // at the same time.
622 unsafe {
623 self.state
624 .query_unchecked_manual_with_ticks(self.world, self.last_run, self.this_run)
625 .iter_many_unique_inner(self.entity_list)
626 .fold(init, func);
627 }
628 }
629 #[cfg(all(not(target_arch = "wasm32"), feature = "multi_threaded"))]
630 {
631 let thread_count = bevy_tasks::ComputeTaskPool::get().thread_num();
632 if thread_count <= 1 {
633 let init = init();
634 // SAFETY: See the safety comment above.
635 unsafe {
636 self.state
637 .query_unchecked_manual_with_ticks(self.world, self.last_run, self.this_run)
638 .iter_many_unique_inner(self.entity_list)
639 .fold(init, func);
640 }
641 } else {
642 // Need a batch size of at least 1.
643 let batch_size = self.get_batch_size(thread_count).max(1);
644 // SAFETY: See the safety comment above.
645 unsafe {
646 self.state.par_many_unique_fold_init_unchecked_manual(
647 init,
648 self.world,
649 &self.entity_list,
650 batch_size,
651 func,
652 self.last_run,
653 self.this_run,
654 );
655 }
656 }
657 }
658 }
659
660 #[cfg(all(not(target_arch = "wasm32"), feature = "multi_threaded"))]
661 fn get_batch_size(&self, thread_count: usize) -> u32 {
662 self.batching_strategy
663 .calc_batch_size(|| self.entity_list.len(), thread_count) as u32
664 }
665}