Repository navigation
perf(sort): optimize coordinate sort with radix sort, memory scaling, and batched I/O - #43
Conversation
… and batched I/O Performance optimizations for coordinate sorting: - Add LSD radix sort for RecordRef (O(n×k) vs O(n log n)) with parallel variant - Scale memory with threads: 768MB/thread by default (was fixed 512MB total) - Add --memory-per-thread flag to control memory scaling behavior - Batch ReadAheadReader: 256 records per channel send (was 1) Bug fixes: - Fix non-stable merge by adding chunk_idx as tie-breaker in heap comparators - Add temp file consolidation (max 64 files, merges oldest when exceeded) Infrastructure: - Add raw BGZF block reader for future direct reading optimization - Add comprehensive radix sort tests Benchmark results (8 threads, 1.9GB BAM, 59M records): - Coordinate sort: 24.7s (was 30.0s) - 18% faster - Outperforms samtools (29.5s) by 16%
📝 WalkthroughWalkthroughAdds per-thread memory scaling for sorting (new 🚥 Pre-merge checks | ✅ 3✅ Passed checks (3 passed)
✏️ Tip: You can configure your own custom pre-merge checks in the settings. ✨ Finishing touches
🧪 Generate unit tests (beta)
Thanks for using CodeRabbit! It's free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out. Comment |
There was a problem hiding this comment.
Actionable comments posted: 4
Caution
Some comments are outside the diff and can’t be posted inline due to platform limitations.
⚠️ Outside diff range comments (1)
src/commands/sort.rs (1)
210-243:⚠️ Potential issue | 🟡 MinorClamp/validate
threadsconsistently withmemory_per_thread.
effective_memoryusesthreads.max(1)but the sorter still gets the rawself.threads. If 0 slips through, logs/memory math diverge and downstream may panic.Proposed fix
- let effective_memory = if self.memory_per_thread { - self.max_memory.saturating_mul(self.threads.max(1)) - } else { - self.max_memory - }; + let threads = self.threads.max(1); + let effective_memory = if self.memory_per_thread { + self.max_memory.saturating_mul(threads) + } else { + self.max_memory + }; @@ - info!("Threads: {}", self.threads); + info!("Threads: {}", threads); @@ - let mut sorter = RawExternalSorter::new(self.order.into()) - .memory_limit(effective_memory) - .threads(self.threads) + let mut sorter = RawExternalSorter::new(self.order.into()) + .memory_limit(effective_memory) + .threads(threads)
🤖 Fix all issues with AI agents
In `@src/lib/sort/raw_bam_reader.rs`:
- Around line 172-209: The docstring for next_record() is incorrect: it claims
the returned Vec<u8> includes the 4-byte block_size prefix but the
implementation in next_record() (which reads block_size and then sets record =
self.decompressed[self.position + 4..self.position + total_size].to_vec())
intentionally strips those 4 bytes and returns only the record payload; update
the doc/comment to state clearly that the returned bytes exclude the 4-byte
block_size prefix (i.e., they are the BAM record payload only) and mention that
total_size = 4 + block_size is used for validation, or alternatively change the
code to include the 4-byte prefix if you prefer keeping the original doc—make
the doc and the behavior of next_record() consistent.
- Around line 124-166: skip_header() calls to ensure_bytes() currently use the ?
operator and ignore the bool return (false on truncated input), leading to
panics when slicing; update skip_header() to check the bool result of
ensure_bytes(...) after each call (for header text, ref name, etc.) and if it is
false return Err(io::Error::new(io::ErrorKind::UnexpectedEof, "truncated BAM
header")); similarly, update read_u32() to treat a false from ensure_bytes(...)
as an UnexpectedEof error instead of proceeding; follow the same pattern used in
next_record() (checking the boolean and returning UnexpectedEof) so all
slices/indexing in skip_header() and read_u32() are guarded against truncated
input.
In `@src/lib/sort/raw.rs`:
- Around line 511-618: The consolidation logic in maybe_consolidate_temp_files
can leave merge_count == 0 and breaks stability by appending the merged file as
newest; fix by computing merge_count = std::cmp::max(1, chunk_files.len()/2) so
we always reduce files when needed, preserve original chunk order by capturing
each file's original index when building readers (use
files_to_merge.iter().enumerate() and store original_chunk_idx in HeapEntry),
use that original_chunk_idx as the tie-breaker in the lt closure instead of
reader_idx, and insert the merged_path back into chunk_files at the position of
the removed oldest files (e.g., chunk_files.insert(0, merged_path) or insert at
the drained range start) rather than push so equal-key records keep the original
relative order.
In `@src/lib/sort/read_ahead.rs`:
- Around line 61-83: with_buffer_size currently always uses BATCH_SIZE for
batch_size which overshoots when buffer_size < BATCH_SIZE; change it to derive
batch_size from buffer_size and guard against zero: compute batch_size =
buffer_size.min(BATCH_SIZE).max(1) and channel_buffer = ((buffer_size +
batch_size - 1) / batch_size).max(1) (ceil division with a 0 guard), then call
Self::with_batch_size(reader, batch_size, channel_buffer); reference symbols:
with_buffer_size, with_batch_size, BATCH_SIZE, channel_buffer, batch_size.
- Fix ensure_bytes() return value handling in raw_bam_reader.rs to return UnexpectedEof errors instead of panicking on truncated input - Fix doc comment for next_record() to accurately describe returned bytes - Fix temp file consolidation stability by inserting merged file at position 0 instead of appending, preserving stable order for equal keys - Add guard for max_temp_files < 2 edge case - Fix with_buffer_size() to handle small buffers correctly using div_ceil and minimum batch size
There was a problem hiding this comment.
Actionable comments posted: 1
🤖 Fix all issues with AI agents
In `@src/lib/sort/read_ahead.rs`:
- Around line 61-75: The receiver must be closed before joining the reader
thread to avoid a drop-time deadlock; change the ReadAheadReader.receiver field
to Option<Receiver<...>> and update constructors such as with_batch_size to
store Some(rx), then modify Drop for ReadAheadReader to take the receiver (e.g.,
self.receiver.take()) and drop it before calling join on the thread handle (use
handle.take() and join()); ensure reader_thread and any code that accesses
receiver account for Option by unwrapping or using take() where appropriate so
the channel is closed prior to join.
🧹 Nitpick comments (2)
src/lib/sort/raw_bam_reader.rs (2)
313-317: Consider validatingbatch_size > 0.A
batch_sizeof 0 causesnext_batch()to always returnNone, appearing as immediate EOF. Consider:assert!(batch_size > 0, "batch_size must be positive");Or use
NonZero<usize>in the signature.
342-349: Test coverage is minimal.Only a constant check exists. When this module is integrated into the sort pipeline, add tests for:
- Valid BAM parsing
- Truncated input handling
- Batch boundary behavior
Acceptable for infrastructure code not yet in use.
Close the channel receiver before joining the reader thread to prevent deadlock when the channel is full and the thread is blocked on send. Wrap receiver in Option to allow taking/dropping it in Drop impl.
Summary
Performance optimizations for coordinate sorting that bring fgumi in line with or exceeding samtools:
RecordRefwith O(n×k) complexity vs O(n log n) comparison sortReadAheadReader(was 1), reducing sync overhead ~256xBug Fixes
chunk_idxas tie-breaker in heap comparatorsInfrastructure
Benchmark Results
Tested on 1.9GB BAM file (59M records), 8 threads:
fgumi now outperforms samtools by 16% on coordinate sorting.
Test plan
fgumi compare bams --mode content