Repository navigation
perf(sort): redesign phase 2 work stealing and bound rayon to --threads - #247
Conversation
|
No actionable comments were generated in the recent review. 🎉 ℹ️ Recent review info⚙️ Run configurationConfiguration used: Organization UI Review profile: CHILL Plan: Pro Run ID: 📒 Files selected for processing (4)
✅ Files skipped from review due to trivial changes (1)
🚧 Files skipped from review as they are similar to previous changes (1)
📝 WalkthroughWalkthroughRefactors Phase 2 into a single 🚥 Pre-merge checks | ✅ 3✅ Passed checks (3 passed)
✏️ Tip: You can configure your own custom pre-merge checks in the settings. ✨ Finishing Touches📝 Generate docstrings
🧪 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: 1
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@src/lib/sort/raw.rs`:
- Around line 1717-1724: The code calls a non-existent method
par_sort_into_chunks on buffer (used when building memory_chunks), which breaks
compilation; either implement par_sort_into_chunks on the concrete buffer types
(RecordBuffer/TemplateRecordBuffer) or revert these branches to use the existing
par_sort() and then materialize chunked Vec<Vec<(RawCoordinateKey, Vec<u8>)>>
via the buffer's current chunking/flush API (i.e., sort with buffer.par_sort()
inside rayon_pool.install and then split the sorted buffer into chunks using the
existing method you already have), and apply the same fix to the other two
occurrences around the referenced blocks.
🪄 Autofix (Beta)
Fix all unresolved CodeRabbit comments on this PR:
- Push a commit to this branch (recommended)
- Create a new PR with the fixes
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: CHILL
Plan: Pro
Run ID: fb5cd1c6-670a-46a4-be8f-5509b71efeac
📒 Files selected for processing (3)
src/lib/sort/raw.rssrc/lib/sort/read_ahead.rssrc/lib/sort/worker_pool.rs
ac35c07 to
22e94bd
Compare
Codecov Report❌ Patch coverage is
Additional details and impacted files@@ Coverage Diff @@
## main #247 +/- ##
==========================================
+ Coverage 89.29% 89.42% +0.12%
==========================================
Files 118 118
Lines 57201 57398 +197
==========================================
+ Hits 51080 51326 +246
+ Misses 6121 6072 -49 ☔ View full report in Codecov by Sentry. 🚀 New features to boost your workflow:
|
There was a problem hiding this comment.
Caution
Some comments are outside the diff and can’t be posted inline due to platform limitations.
⚠️ Outside diff range comments (1)
src/lib/sort/raw.rs (1)
2451-2460:⚠️ Potential issue | 🟠 MajorMake Phase 2 teardown run on every exit path.
From
pool.set_phase(phase::PHASE2)until the manual cleanup block, several?sites can return early. That leaves the pool inPHASE2and keepsphase2_filespublished on merge errors. Wrap this in a drop guard so the reset always happens.Suggested shape
+struct Phase2Cleanup<'a> { + pool: &'a Arc<SortWorkerPool>, + active: bool, +} + +impl Drop for Phase2Cleanup<'_> { + fn drop(&mut self) { + if self.active { + self.pool.set_phase(phase::LEGACY); + self.pool.clear_phase2_files(); + } + } +} + let mut consumer: Option<MainThreadChunkConsumer<K>> = if num_disk > 0 { let files = pool.phase2_files(); let consumer = MainThreadChunkConsumer::new( files, pool.decompress_error_flag(), pool.chunk_read_error_flag(), pool.worker_panicked_flag(), ); pool.set_phase(phase::PHASE2); Some(consumer) } else { None }; + let _phase2_cleanup = consumer.as_ref().map(|_| Phase2Cleanup { pool, active: true }); ... - if consumer.is_some() { - pool.set_phase(phase::LEGACY); - drop(consumer.take()); - pool.clear_phase2_files(); - }Also applies to: 2472-2475, 2531-2558
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@src/lib/sort/raw.rs` around lines 2451 - 2460, The code sets pool.set_phase(phase::PHASE2) and then constructs a MainThreadChunkConsumer and performs operations that may return early via ?; if an early return occurs the pool remains in PHASE2 and phase2_files stay published. Fix by introducing a RAII drop guard (e.g., a small struct ResetPhaseOnDrop) created immediately after calling pool.set_phase(phase::PHASE2) that in its Drop impl resets the phase and clears/hides phase2_files, then ensure the guard is dropped only after the manual cleanup block completes (or explicitly into_inner when success) so the reset runs on every exit path; apply the same pattern around the other regions mentioned (near lines where MainThreadChunkConsumer is created and where phase is set, e.g., the blocks covering 2472-2475 and 2531-2558).
♻️ Duplicate comments (1)
src/lib/sort/raw.rs (1)
1721-1722:⚠️ Potential issue | 🔴 CriticalReplace
par_sort_into_chunksor add that API first.These call sites still target a method that
RecordBuffer/TemplateRecordBufferdo not expose, so this is a hard compile break. Either addpar_sort_into_chunkson those buffer types or keeppar_sort()and materialize merge chunks with the existing iterators/split logic.Also applies to: 1893-1894, 2298-2299
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@src/lib/sort/raw.rs` around lines 1721 - 1722, Call sites use buffer.par_sort_into_chunks which RecordBuffer/TemplateRecordBuffer lack, causing compile errors; either implement par_sort_into_chunks on those types or revert the call sites to use the existing par_sort() and then materialize chunks via the current iterators/split logic. Update occurrences such as the timer.time_sort(|| rayon_pool.install(|| buffer.par_sort_into_chunks(self.threads))) (and the similar sites around the other mentioned lines) to either call a newly added RecordBuffer::par_sort_into_chunks/TemplateRecordBuffer::par_sort_into_chunks that returns chunked iterators compatible with the merge stage, or change them back to buffer.par_sort() and follow with the code that splits the sorted buffer into merge chunks using the existing iterator/split helpers so callers compile without the new API.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Outside diff comments:
In `@src/lib/sort/raw.rs`:
- Around line 2451-2460: The code sets pool.set_phase(phase::PHASE2) and then
constructs a MainThreadChunkConsumer and performs operations that may return
early via ?; if an early return occurs the pool remains in PHASE2 and
phase2_files stay published. Fix by introducing a RAII drop guard (e.g., a small
struct ResetPhaseOnDrop) created immediately after calling
pool.set_phase(phase::PHASE2) that in its Drop impl resets the phase and
clears/hides phase2_files, then ensure the guard is dropped only after the
manual cleanup block completes (or explicitly into_inner when success) so the
reset runs on every exit path; apply the same pattern around the other regions
mentioned (near lines where MainThreadChunkConsumer is created and where phase
is set, e.g., the blocks covering 2472-2475 and 2531-2558).
---
Duplicate comments:
In `@src/lib/sort/raw.rs`:
- Around line 1721-1722: Call sites use buffer.par_sort_into_chunks which
RecordBuffer/TemplateRecordBuffer lack, causing compile errors; either implement
par_sort_into_chunks on those types or revert the call sites to use the existing
par_sort() and then materialize chunks via the current iterators/split logic.
Update occurrences such as the timer.time_sort(|| rayon_pool.install(||
buffer.par_sort_into_chunks(self.threads))) (and the similar sites around the
other mentioned lines) to either call a newly added
RecordBuffer::par_sort_into_chunks/TemplateRecordBuffer::par_sort_into_chunks
that returns chunked iterators compatible with the merge stage, or change them
back to buffer.par_sort() and follow with the code that splits the sorted buffer
into merge chunks using the existing iterator/split helpers so callers compile
without the new API.
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: CHILL
Plan: Pro
Run ID: 0978eb85-d5dd-4302-bec0-b19eedce75e5
📒 Files selected for processing (3)
src/lib/sort/raw.rssrc/lib/sort/read_ahead.rssrc/lib/sort/worker_pool.rs
✅ Files skipped from review due to trivial changes (1)
- src/lib/sort/worker_pool.rs
22e94bd to
50071ce
Compare
|
@coderabbitai review |
✅ Actions performedReview triggered.
|
There was a problem hiding this comment.
Actionable comments posted: 1
🧹 Nitpick comments (1)
src/lib/sort/inline_buffer.rs (1)
248-294: Factor the chunking skeleton once.These two methods now duplicate the same threshold check, chunk sizing, parallel chunk sort, and materialization flow. A small internal helper would keep future tuning and fixes in one place.
Also applies to: 844-876
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@src/lib/sort/inline_buffer.rs` around lines 248 - 294, Extract the duplicated chunking/sorting/materialization logic into a private helper (e.g., fn materialize_sorted_chunks(&mut self, threads: usize) -> Vec<Vec<(RawCoordinateKey, Vec<u8>)>>) that encapsulates the threshold checks (using RADIX_THRESHOLD, n = self.refs.len(), and the threads <= 1 / small-n single-threaded path), chunk_size calculation (n.div_ceil(threads)), parallel radix-sorting of chunks via self.refs.par_chunks_mut(chunk_size) with radix_sort_record_refs, and the final materialization using self.refs.chunks(chunk_size) mapping each Ref to (RawCoordinateKey { sort_key: r.sort_key }, self.get_record(r).to_vec()). Replace the bodies of par_sort_into_chunks and the other similar method (the one at lines ~844-876) to call this new helper so tuning and fixes live in one place; keep all semantics and use of radix_sort_record_refs, get_record, refs, and RADIX_THRESHOLD unchanged.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@src/lib/sort/raw.rs`:
- Around line 2136-2158: The emitted chunk boundaries currently don't match the
ranges sorted by par_chunks_mut because you split from the tail with split_off;
change the emission logic so chunks are carved from the head in the same
fixed-size ranges used by par_chunks_mut (i.e., [0..chunk_size),
[chunk_size..2*chunk_size), ...), ensuring each pushed Vec corresponds to a
sorted sub-slice; specifically adjust the code that builds chunks (the
remaining/split_off loop that produces chunks) to iterate over entries from the
start in chunk_size increments (handling the final short chunk) so that
chunk_size, par_chunks_mut, remaining, split_off and chunks align with the k-way
merge assumption that every source chunk is sorted.
---
Nitpick comments:
In `@src/lib/sort/inline_buffer.rs`:
- Around line 248-294: Extract the duplicated chunking/sorting/materialization
logic into a private helper (e.g., fn materialize_sorted_chunks(&mut self,
threads: usize) -> Vec<Vec<(RawCoordinateKey, Vec<u8>)>>) that encapsulates the
threshold checks (using RADIX_THRESHOLD, n = self.refs.len(), and the threads <=
1 / small-n single-threaded path), chunk_size calculation (n.div_ceil(threads)),
parallel radix-sorting of chunks via self.refs.par_chunks_mut(chunk_size) with
radix_sort_record_refs, and the final materialization using
self.refs.chunks(chunk_size) mapping each Ref to (RawCoordinateKey { sort_key:
r.sort_key }, self.get_record(r).to_vec()). Replace the bodies of
par_sort_into_chunks and the other similar method (the one at lines ~844-876) to
call this new helper so tuning and fixes live in one place; keep all semantics
and use of radix_sort_record_refs, get_record, refs, and RADIX_THRESHOLD
unchanged.
🪄 Autofix (Beta)
Fix all unresolved CodeRabbit comments on this PR:
- Push a commit to this branch (recommended)
- Create a new PR with the fixes
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: CHILL
Plan: Pro
Run ID: 988417c3-a276-4e96-8096-c2d89d4d6942
📒 Files selected for processing (4)
src/lib/sort/inline_buffer.rssrc/lib/sort/raw.rssrc/lib/sort/read_ahead.rssrc/lib/sort/worker_pool.rs
🚧 Files skipped from review as they are similar to previous changes (1)
- src/lib/sort/read_ahead.rs
Two independent problems were hurting the N+2 sort worker-pool path
on large inputs; this squash fixes both and lands a post-review
cleanup sweep.
## Phase 2 redesign: per-file work-stealing state
Phase 2 regressed badly on WES (1kg-wes-HG00100.bam, 13G, 219M
records): 450s+ wall and ~42 GB peak RSS under --threads 8
--max-memory 768M, driven by an over-coupled producer/consumer with
global queues that deadlocked under backpressure and left workers
idle while OOMing.
Replace it with per-file state:
struct Phase2FileState {
reader: Mutex<Phase2Reader>, // raw BGZF block reads
reader_eof: AtomicBool, // hot-path fast check
raw_blocks: Mutex<VecDeque<...>>, // bounded FIFO
decompressed: Mutex<ReorderBuffer>, // in-order decompressed
decomp_in_flight: AtomicUsize, // race-safe drain
}
Workers scan files with try_lock everywhere and steal work across
spill files. Lock order is strict: reader -> raw_blocks ->
decompressed.
- Deadlock-free admission: a popped raw block is always admitted when
it is the gap-filler the reorder buffer is waiting for, even if the
decompressed cap is full, preventing head-of-line deadlock where
every worker holds a stale raw block while the buffer waits on a
missing serial.
- Race fix: is_drained() must observe decomp_in_flight > 0 to avoid
exiting Phase 2 between a worker popping a raw block and inserting
the decompressed bytes.
- FIFO bug fix: serial assignment and the matching push into
raw_blocks happen under both the reader and raw_blocks locks held
simultaneously. Releasing the reader lock first allowed another
worker to interleave a higher serial ahead of the lower one and
break the gap-filler invariant.
- is_drained() fast-paths on a lock-free reader_eof atomic instead of
acquiring the reader mutex; the common case (reader not at EOF)
now exits after a single atomic load, versus three mutexes per
decompressed block on the original code.
## Rayon oversubscription
Rayon's global pool defaults to num_cpus::get(), which silently
violated --threads on machines with more physical cores than the
requested thread count. Every par_sort / par_sort_into_chunks /
par_sort_unstable_by / par_chunks_mut call in the raw sort path now
runs under a rayon::ThreadPool built with num_threads = self.threads,
via rayon_pool.install(|| ...). rayon::current_num_threads() inside
parallel_radix_sort_* now returns the requested thread count, so
chunk sizing is consistent with the worker pool sizing.
Oversubscription with the SortWorkerPool is not introduced: every
call site is preceded by drain_pending_spill, which joins the prior
chunk's I/O thread, so the sort worker pool is provably idle when
rayon fans out.
## Additional fixes
- Use par_sort in the three in-memory sort fallbacks (coordinate,
coordinate-index, template-coordinate). The single-threaded
buffer.sort() serialized the sort step onto one core even with
--threads set high when the entire input fit in memory.
- In-memory split path in RawExternalSorter replaces the O(n^2)
drain(..chunk_size) loop with O(n) split_off from the tail. The
old loop memmoved the tail on every drain and became a measurable
cost once entries grew into the millions under tight --max-memory
budgets. The merge consumer accepts chunks in any order so the
original layout need not be preserved.
- Drop Phase2FileState::source_id: it was always equal to the file's
index in the shared vector. Log lines use the cursor index and
set_phase2_files simplifies to &[PathBuf].
- Drop SourceParserState::eof and the unused _worker parameter on
is_phase_complete; both were written or threaded through but never
read.
- Refresh stale module, struct, and function docs to describe the
per-file work-stealing Phase 2 model. No stale references to
ReadChunkBlocks, decompressed_chunks, or GenericKeyedChunkReader
remain.
## Result
WES (1kg-wes-HG00100.bam, --threads 8 --max-memory 4g):
before (5df7e17): 450s+ wall, ~42 GB peak RSS (OOM-prone)
after: 139 s wall, 37.6 GB peak RSS, 0 sort violations
50071ce to
789fc66
Compare
Summary
Fixes two independent problems hurting the N+2 sort worker-pool path on large inputs, plus a post-review cleanup sweep.
Phase 2 redesign: per-file work-stealing state
Phase 2 regressed badly on WES (1kg-wes-HG00100.bam, 13 GB, 219M records): 450s+ wall and ~42 GB peak RSS under
--threads 8 --max-memory 768M, driven by an over-coupled producer/consumer with global queues that deadlocked under backpressure and left workers idle while OOMing.Replace it with per-file state:
Workers scan files with
try_lockeverywhere and steal work across spill files. Lock order is strict:reader -> raw_blocks -> decompressed.is_drained()must observedecomp_in_flight > 0to avoid exiting Phase 2 between a worker popping a raw block and inserting the decompressed bytes.raw_blockshappen under both the reader and raw_blocks locks held simultaneously. Releasing the reader lock first allowed another worker to interleave a higher serial ahead of the lower one and break the gap-filler invariant.is_drained()fast-paths on a lock-freereader_eofatomic; the common case (reader not at EOF) exits after a single atomic load, versus three mutex acquisitions per decompressed block.Rayon oversubscription
Rayon's global pool defaults to
num_cpus::get(), which silently violated--threadson machines with more physical cores than the requested thread count. Everypar_sort/par_sort_into_chunks/par_sort_unstable_by/par_chunks_mutcall in the raw sort path now runs under arayon::ThreadPoolbuilt withnum_threads = self.threads, viarayon_pool.install(|| ...).rayon::current_num_threads()insideparallel_radix_sort_*now returns the requested thread count, so chunk sizing is consistent with the worker-pool sizing.Oversubscription with the
SortWorkerPoolis not introduced: every call site is preceded bydrain_pending_spill, which joins the prior chunk's I/O thread, so the sort worker pool is provably idle when rayon fans out.Additional fixes
par_sortin the three in-memory sort fallbacks (coordinate, coordinate-index, template-coordinate). The single-threadedbuffer.sort()serialized the sort step onto one core even with--threadsset high when the entire input fit in memory.drain(..chunk_size)loop with O(n)split_offfrom the tail. The old loop memmoved the tail on every drain and became a measurable cost onceentriesgrew into the millions under tight--max-memorybudgets.Phase2FileState::source_id(always equal to the file's index in the shared vector),SourceParserState::eof, and the unused_workerparameter onis_phase_complete— all written or threaded through but never read.Result
WES (
1kg-wes-HG00100.bam,--threads 8 --max-memory 4g):Test plan