rust
40 lines · 8 steps
A parallel map-reduce with scoped threads
Splitting a slice across worker threads, mapping and folding each chunk in parallel, then combining the partial results.
Explained by
highlit
1use std::thread;
2
3pub fn parallel_map_reduce<T, M, R, F, G, H>(
4 data: &[T],
5 workers: usize,
6 map: F,
7 reduce: G,
8 identity: H,
9) -> R
10where
11 T: Sync,
12 M: Send,
13 R: Send,
14 F: Fn(&T) -> M + Sync,
15 G: Fn(R, M) -> R + Sync,
16 H: Fn() -> R + Sync,
17{
18 if data.is_empty() {
19 return identity();
20 }
21 let workers = workers.max(1).min(data.len());
22 let chunk_size = (data.len() + workers - 1) / workers;
23
24 let partials: Vec<R> = thread::scope(|scope| {
25 let handles: Vec<_> = data
26 .chunks(chunk_size)
27 .map(|chunk| {
28 scope.spawn(|| {
29 chunk
30 .iter()
31 .map(&map)
32 .fold(identity(), &reduce)
33 })
34 })
35 .collect();
36 handles.into_iter().map(|h| h.join().unwrap()).collect()
37 });
38
39 partials.into_iter().fold(identity(), &reduce)
40}
01 / 01
STEP 01
‹ swipe to step through ›
Walkthrough
Space play
←→ step
click any line
Three takeaways
- 1Scoped threads let closures borrow stack data directly, avoiding clones or Arc just to share the input slice.
- 2An identity-producing closure lets each chunk fold independently, so partials can be combined the same way at the end.
- 3Sync and Send bounds encode exactly which values cross thread boundaries, making the parallelism safe at compile time.
Related explainers
rust
use serde::Deserialize; #[derive(Debug, Deserialize)] #[serde(untagged)]
Parsing flexible JSON shapes with serde
deserialization
enums
json
Intermediate
6 steps
rust
use std::cmp::{Ordering, Reverse}; use std::collections::BinaryHeap; use std::fs::File; use std::io::{self, BufRead, BufReader, BufWriter, Lines, Write};
K-way merge of sorted logs in Rust
binary-heap
k-way-merge
streaming-io
Intermediate
8 steps
rust
use axum::{ extract::{Path, State}, response::sse::{Event, KeepAlive, Sse}, };
Streaming import progress with SSE in Axum
server-sent-events
streams
watch-channel
Advanced
7 steps
rust
use std::f64::consts::PI; #[derive(Debug, Clone, Copy)] pub struct LatLng {
Building geographic bounding boxes in Rust
geospatial
value-types
option
Intermediate
7 steps
go
func (w *Watcher) resetDebounce(d time.Duration) { if !w.timer.Stop() { select { case <-w.timer.C:
Debouncing a stream of events in Go
debounce
timers
channels
Advanced
7 steps
rust
use chrono::{Duration, NaiveDate}; #[derive(Debug)] pub struct DateRange {
Parsing and iterating date ranges in Rust
error-handling
iterators
parsing
Intermediate
7 steps
Share this explainer
Here's the card — post it anywhere.
Made with highlit — turn any snippet into a walkthrough like this in about a minute.
Explain your code
Embed this explainer
Drop the interactive walkthrough into a blog or docs. Views never cost a credit.
<iframe src="https://highlit.co/explainers/a-parallel-map-reduce-with-scoped-threads-explained-rust-6404/embed?autoplay=1" width="100%" height="520" loading="lazy" style="border:0"></iframe>
Autoplay is on by default — add ?autoplay=0 to start paused.