Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 9 additions & 4 deletions parallel_bzip2_decoder/benches/scanner_benchmark.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -53,15 +54,17 @@ 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(
BenchmarkId::from_parameter(format!("{}MB", size_mb)),
&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;
Expand Down Expand Up @@ -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;
Expand Down
4 changes: 2 additions & 2 deletions parallel_bzip2_decoder/src/decoder.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<dyn AsRef<[u8]> + 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
Expand Down
17 changes: 8 additions & 9 deletions parallel_bzip2_decoder/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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.
///
Expand Down Expand Up @@ -140,26 +141,24 @@ 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<T>(data: Arc<T>) -> 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
// This maintains cache locality and limits memory usage
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
Expand Down Expand Up @@ -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));
}
});
Expand Down