From 7ee5eec9042608ef871c2b8240aae164cadb3ae1 Mon Sep 17 00:00:00 2001 From: Jan Ole Zabel Date: Tue, 8 Sep 2026 22:37:56 +0200 Subject: [PATCH 1/2] Remove unneeded full data copy --- parallel_bzip2_decoder/src/decoder.rs | 4 ++-- parallel_bzip2_decoder/src/lib.rs | 17 ++++++++--------- 2 files changed, 10 insertions(+), 11 deletions(-) diff --git a/parallel_bzip2_decoder/src/decoder.rs b/parallel_bzip2_decoder/src/decoder.rs index 3288712..1023422 100644 --- a/parallel_bzip2_decoder/src/decoder.rs +++ b/parallel_bzip2_decoder/src/decoder.rs @@ -147,14 +147,14 @@ impl Bz2Decoder { // Channel for sending decompressed blocks back to the reader // Sized at 2x thread count to allow some buffering without excessive memory use let (result_sender, result_receiver) = bounded(rayon::current_num_threads() * 2); - let data_ref: Arc + Send + Sync> = data; + let data_ref = data; let data_clone = data_ref.clone(); // Spawn the driver thread that coordinates scanning and decompression std::thread::spawn(move || { let slice = data_clone.as_ref().as_ref(); // Get block boundaries from the scanner - let task_receiver = scan_blocks(slice); + let task_receiver = scan_blocks(data_clone.clone()); // Parallel decompression using Rayon // par_bridge() allows us to process an iterator in parallel diff --git a/parallel_bzip2_decoder/src/lib.rs b/parallel_bzip2_decoder/src/lib.rs index 0d3f93f..0d2ad1b 100644 --- a/parallel_bzip2_decoder/src/lib.rs +++ b/parallel_bzip2_decoder/src/lib.rs @@ -100,6 +100,7 @@ use bzip2::read::BzDecoder; use crossbeam_channel::bounded; use std::collections::HashMap; use std::io::Read; +use std::sync::Arc; /// Scans bzip2 data for block boundaries and returns them via a channel. /// @@ -140,16 +141,14 @@ use std::io::Read; /// println!("Block from bit {} to bit {}", start, end); /// } /// ``` -pub fn scan_blocks(data: &[u8]) -> crossbeam_channel::Receiver<(u64, u64)> { +pub fn scan_blocks(data: Arc) -> crossbeam_channel::Receiver<(u64, u64)> +where + T: AsRef<[u8]> + Send + Sync + 'static + ?Sized, +{ // Channel for sending block boundaries to the caller // Buffer size of 100 allows good throughput without excessive memory use let (task_sender, task_receiver) = bounded(100); - // Clone data into an Arc for safe sharing across threads - let data_vec = data.to_vec(); - let data_arc = std::sync::Arc::new(data_vec); - let data_clone = data_arc.clone(); - std::thread::spawn(move || { let scanner = Scanner::new(); // Small buffer for chunks to prevent scanning too far ahead @@ -157,9 +156,9 @@ pub fn scan_blocks(data: &[u8]) -> crossbeam_channel::Receiver<(u64, u64)> { let (chunk_tx, chunk_rx) = bounded(4); // Spawn the actual scanning in a background thread - let scan_data = data_clone.clone(); + let scan_data = data.clone(); let _scan_handle = std::thread::spawn(move || { - scanner.scan_stream(&scan_data, 0, chunk_tx); + scanner.scan_stream(scan_data.as_ref().as_ref(), 0, chunk_tx); }); // Reorder chunks and convert markers to block boundaries @@ -200,7 +199,7 @@ pub fn scan_blocks(data: &[u8]) -> crossbeam_channel::Receiver<(u64, u64)> { // Handle edge case: block without EOS marker (truncated file) if let Some(start) = current_block_start { - let end = (data_clone.len() as u64) * 8; + let end = (data.as_ref().as_ref().len() as u64) * 8; let _ = task_sender.send((start, end)); } }); From 08e38fcabe465d96990ab735aac2c703b44f6e26 Mon Sep 17 00:00:00 2001 From: Jan Ole Zabel Date: Tue, 8 Sep 2026 23:40:44 +0200 Subject: [PATCH 2/2] Update benchmark code --- parallel_bzip2_decoder/benches/scanner_benchmark.rs | 13 +++++++++---- 1 file changed, 9 insertions(+), 4 deletions(-) diff --git a/parallel_bzip2_decoder/benches/scanner_benchmark.rs b/parallel_bzip2_decoder/benches/scanner_benchmark.rs index 074d54a..6e47f7e 100644 --- a/parallel_bzip2_decoder/benches/scanner_benchmark.rs +++ b/parallel_bzip2_decoder/benches/scanner_benchmark.rs @@ -5,6 +5,7 @@ use pprof::criterion::{Output, PProfProfiler}; use std::fs; use std::path::Path; use std::process::Command; +use std::sync::Arc; /// Generate a test file of the specified size in MB fn generate_test_file(size_mb: usize) -> String { @@ -53,7 +54,9 @@ fn bench_scanner(c: &mut Criterion) { for size_mb in [1, 10, 50].iter() { let bz2_file = generate_test_file(*size_mb); - let data = std::fs::read(&bz2_file).expect("Failed to read test file"); + let data: Arc<[u8]> = std::fs::read(&bz2_file) + .expect("Failed to read test file") + .into(); group.throughput(Throughput::Bytes(data.len() as u64)); group.bench_with_input( @@ -61,7 +64,7 @@ fn bench_scanner(c: &mut Criterion) { &data, |b, data| { b.iter(|| { - let receiver = scan_blocks(data); + let receiver = scan_blocks(data.clone()); let mut count = 0; while receiver.recv().is_ok() { count += 1; @@ -118,14 +121,16 @@ fn bench_scanner_multistream(c: &mut Criterion) { return; } - let data = std::fs::read(&bz2_filename).expect("Failed to read test file"); + let data: Arc<[u8]> = std::fs::read(&bz2_filename) + .expect("Failed to read test file") + .into(); let mut group = c.benchmark_group("scanner_multistream"); group.throughput(Throughput::Bytes(data.len() as u64)); group.bench_function("scan_multistream", |b| { b.iter(|| { - let receiver = scan_blocks(&data); + let receiver = scan_blocks(data.clone()); let mut count = 0; while receiver.recv().is_ok() { count += 1;