parallelism

This commit is contained in:
timeshifter
2026-07-09 19:12:38 +02:00
parent 3080d8a64c
commit b9e05f3cdf
2 changed files with 86 additions and 31 deletions
+1
View File
@@ -6,3 +6,4 @@
71 mmap, reduced UTF8 parsing (names) 71 mmap, reduced UTF8 parsing (names)
50 manual number parsing 50 manual number parsing
48 unchecked unwrap 48 unchecked unwrap
8 12-thread parallel
+87 -33
View File
@@ -2,50 +2,51 @@ use fxhash::FxHashMap;
use memmap2::MmapOptions; use memmap2::MmapOptions;
fn main() { fn main() {
let mut data = FxHashMap::default();
let file = std::fs::File::open("measurements.txt").unwrap(); 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 // SAFETY: we are not going to modify the file while the program is running. No lock required
// for now. // for now.
let mut datas = vec![];
let mmap = unsafe { MmapOptions::new().map(&file) }.unwrap(); let mmap = unsafe { MmapOptions::new().map(&file) }.unwrap();
mmap.advise(memmap2::Advice::Sequential).unwrap(); mmap.advise(memmap2::Advice::Sequential).unwrap();
let mut pointer = 0; let length = mmap.len();
let mut length = 0;
loop { let nthreads = std::thread::available_parallelism().unwrap();
length = 0; let (tx, rx) = std::sync::mpsc::sync_channel(nthreads.into());
while let Some(ch) = &mmap.get(pointer + length) { std::thread::scope(|scope| {
if **ch == b'\n' { let mut handles = vec![];
break;
} else { let chunk_size = length / nthreads;
length += 1; 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);
for h in handles {
h.join().unwrap().unwrap();
}
});
for data in rx.iter() {
datas.push(data);
} }
if length == 0 || pointer + length == mmap.len() { let data: FxHashMap<&[u8], (i64, i64, usize, i64)> =
break; datas.into_iter().flat_map(|d| d.into_iter()).collect();
}
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;
}
let mut stations: Vec<_> = data.keys().collect(); let mut stations: Vec<_> = data.keys().collect();
stations.sort_unstable(); stations.sort_unstable();
@@ -71,6 +72,59 @@ fn main() {
println!("}}"); 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 { fn parse_number(number: &[u8]) -> i64 {
unsafe { unsafe {
let sign = if *number.first().unwrap_unchecked() == b'-' { let sign = if *number.first().unwrap_unchecked() == b'-' {