From b9e05f3cdfa464af67310252c2b9469d673c2529 Mon Sep 17 00:00:00 2001 From: timeshifter Date: Thu, 9 Jul 2026 19:12:38 +0200 Subject: [PATCH] parallelism --- benchmark.txt | 1 + src/main.rs | 116 ++++++++++++++++++++++++++++++++++++-------------- 2 files changed, 86 insertions(+), 31 deletions(-) diff --git a/benchmark.txt b/benchmark.txt index 9db7c6f..dd39ac7 100644 --- a/benchmark.txt +++ b/benchmark.txt @@ -6,3 +6,4 @@ 71 mmap, reduced UTF8 parsing (names) 50 manual number parsing 48 unchecked unwrap + 8 12-thread parallel diff --git a/src/main.rs b/src/main.rs index dccdf77..c6fff00 100644 --- a/src/main.rs +++ b/src/main.rs @@ -2,51 +2,52 @@ use fxhash::FxHashMap; use memmap2::MmapOptions; fn main() { - let mut data = FxHashMap::default(); - let file = std::fs::File::open("measurements.txt").unwrap(); // SAFETY: we are not going to modify the file while the program is running. No lock required // for now. + + let mut datas = vec![]; + let mmap = unsafe { MmapOptions::new().map(&file) }.unwrap(); mmap.advise(memmap2::Advice::Sequential).unwrap(); - let mut pointer = 0; - let mut length = 0; + let length = mmap.len(); - loop { - length = 0; + let nthreads = std::thread::available_parallelism().unwrap(); + let (tx, rx) = std::sync::mpsc::sync_channel(nthreads.into()); - while let Some(ch) = &mmap.get(pointer + length) { - if **ch == b'\n' { - break; - } else { - length += 1; - } + std::thread::scope(|scope| { + let mut handles = vec![]; + + let chunk_size = length / nthreads; + let mut start = 0; + + for _ in 0..nthreads.into() { + let end = (start + chunk_size).min(length); + let end = find_newline_idx(&mmap, end); + + let map = &mmap[start..end]; + + let this_tx = tx.clone(); + let handle = scope.spawn(move || this_tx.send(single_thread(map))); + handles.push(handle); + + start = end; } + drop(tx); - if length == 0 || pointer + length == mmap.len() { - break; + for h in handles { + h.join().unwrap().unwrap(); } - let line = &mmap[pointer..pointer + length]; + }); - let mut split = line.split(|&c| c == b';'); - let station = unsafe { split.next().unwrap_unchecked() }; - let number = unsafe { split.next().unwrap_unchecked() }; - let parsed = parse_number(number); - // panic!(); - - data.entry(station) - .and_modify(|values: &mut (i64, i64, usize, i64)| { - values.0 = values.0.min(parsed); - values.1 += parsed; - values.2 += 1; - values.3 = values.3.max(parsed); - }) - .or_insert((parsed, parsed, 1_usize, parsed)); - - pointer += length + 1; + for data in rx.iter() { + datas.push(data); } + let data: FxHashMap<&[u8], (i64, i64, usize, i64)> = + datas.into_iter().flat_map(|d| d.into_iter()).collect(); + let mut stations: Vec<_> = data.keys().collect(); stations.sort_unstable(); @@ -71,6 +72,59 @@ fn main() { println!("}}"); } +fn find_newline_idx(mmap: &memmap2::Mmap, start: usize) -> usize { + let mut idx = start; + while let Some(c) = mmap.get(idx) { + if *c == b'\n' { + return idx; + } + idx += 1; + } + mmap.len() +} + +fn single_thread(mmap: &[u8]) -> FxHashMap<&[u8], (i64, i64, usize, i64)> { + let mut data = FxHashMap::default(); + + let mut pointer = 0; + let mut length = 0; + + loop { + length = 0; + + while let Some(ch) = &mmap.get(pointer + length) { + if **ch == b'\n' { + break; + } else { + length += 1; + } + } + + if pointer + length == mmap.len() { + break; + } + let line = &mmap[pointer..pointer + length]; + + let mut split = line.split(|&c| c == b';'); + let station = unsafe { split.next().unwrap_unchecked() }; + let number = unsafe { split.next().unwrap_unchecked() }; + let parsed = parse_number(number); + + data.entry(station) + .and_modify(|values: &mut (i64, i64, usize, i64)| { + values.0 = values.0.min(parsed); + values.1 += parsed; + values.2 += 1; + values.3 = values.3.max(parsed); + }) + .or_insert((parsed, parsed, 1_usize, parsed)); + + pointer += length + 1; + } + + data +} + fn parse_number(number: &[u8]) -> i64 { unsafe { let sign = if *number.first().unwrap_unchecked() == b'-' {