Skip to content

perf: parallel gather for the in-memory sort fast path - #966

Merged
nh13 merged 1 commit into
mainfrom
nh/sort-merge-parallel-gather
Sep 17, 2026
Merged

nh13 merged 1 commit into
mainfrom
nh/sort-merge-parallel-gather

Conversation

@nh13

@nh13 nh13 commented Sep 14, 2026 •

Copy link
Copy Markdown
Member

What

fgumi sort's in-memory fast path gathered the whole sorted chunk into output blocks serially on one detached coordination thread — a flat ~1.2 s cost that did not scale with --threads and floored the sort's wall clock while the worker pool sat mostly idle. This PR fans that gather across a bounded, step-owned rayon pool sized to the phase-2 (merge) thread budget.

Only the fast path (a single already-sorted in-memory chunk, no k-way merge) is touched. The spill / k-way-merge path is deliberately untouched — it is a genuinely serial loser-tree walk whose winner order cannot be produced out of order.

How

  1. Plan block boundaries up front, off the cold arena. New record_len accessors (InMemoryChunk / TemplateMemChunk / MemoryChunkErased) read only the per-record len index. plan_fast_path_blocks reproduces the serial count/byte-cap split byte-for-byte from lengths alone. A MergeBatchBuilder::FRAME_OVERHEAD_PER_RECORD const (4 for the framed BlockBuilder, 0 for RecordBatchBuilder) keeps the byte-cap plan in lockstep with each builder's total_bytes() growth.
  2. Gather windows in parallel, emit in strict ascending dense ordinal order. SortMerge is Detached, so its output edge has no reorder stage (the by-ordinal reassembly is downstream at BgzfCompress and needs a dense, gap-free stream). Gathered-but-unpushed blocks ride the existing held-slot backpressure idiom plus a small fast_pending queue, drained in order at the top of the next dispatch.
  3. Bounded + gated. The window is fast_path_threads * 4 blocks. Gated on FAST_PATH_PARALLEL_MIN_RECORDS = 64Ki so small chunks keep the cheaper serial gather; a degenerate single-block plan falls through to serial rather than spinning up a pool. The pool is step-owned and bounded to num_phase2_threads so it never oversubscribes past --threads.

Results (EC2 c8g.8xlarge, 32 vCPU Graviton4, coordinate sort, 20 M records)

threads wall before wall after speedup
16 1.75 s 0.92 s ~1.9×
32 1.84 s 0.93 s ~2.0×

Total CPU-seconds is roughly flat — the win is converting serial gather time into parallel gather time. 16 threads is the throughput sweet spot (21.9 M rec/s).

Correctness

Output is byte-for-byte identical to the serial gather. Verified with fgumi compare bams --command sort (record-content + order) across coordinate / template-coordinate / queryname at 16 threads, and against the spill regime (fast path not taken) — all IDENTICAL. Parity unit tests assert the parallel gather emits identical dense ordinals and framed bytes as the serial gather across count-cap, byte-cap, small-multi-block, partial-tail, and a deterministic backpressure regime (small output queue forcing sustained mid-window reject/resume) for both the BlockOutput and RecordBatchOutput builders.

Gates

cargo ci-fmt, cargo ci-lint, cargo ci-doc clean; cargo ci-test — 10219 passed, 31 skipped.

Risk: sort output stays byte-identical through serial-equivalent block planning and ordered emission; unsafe code: none added, and the CLAUDE.md allowlist remains unchanged; memory and backpressure policy: bounded parallel gathering uses the merge-thread budget.

Fix: parallelize only large, multi-block in-memory gathers. Keep small inputs serial. Keep spill and k-way merge serial.

  • Add record_len accessors for in-memory chunks.
  • Add per-record framing overhead to MergeBatchBuilder.
  • Add configurable fast-path thread limits.
  • Pass phase-2 thread limits from pipeline builders.
  • Add parity tests for bytes, framing, ordinals, boundaries, and backpressure.
  • Reported tests and benchmarks pass, with approximately 1.9–2.0× speedup on the stated large coordinate-sort benchmark.

The streaming sort's SortMerge fast path (single already-sorted in-memory
chunk, no k-way merge) gathered every record into output blocks serially on
the lone detached coordination thread. Profiling a 20M-record coordinate sort
at 32 threads showed this gather as a flat ~1.2s serial cost that does not
scale with thread count and floors the wall clock while the worker pool sits
~80% idle -- its CPU is per-record record_bytes() arena slices plus the
extend_from_slice memcpy into block buffers.

Fan that gather across a bounded, step-owned rayon pool sized to the phase-2
(merge) thread budget:

- Precompute output-block boundaries up front from the per-record lengths
  alone (new InMemoryChunk/MemoryChunkErased::record_len, off the cold arena),
  reproducing the serial count/byte-cap split byte-for-byte. A
  MergeBatchBuilder::FRAME_OVERHEAD_PER_RECORD const keeps the byte-cap plan in
  lockstep with each output builder (BlockOutput frames +4/record,
  RecordBatchOutput +0).
- Gather a bounded window of blocks in parallel (chunk is Send+Sync;
  record_bytes is a read-only slice), then push them in strict ascending dense
  ordinal order through the existing held-slot backpressure idiom -- required,
  since the Detached output edge has no reorder stage (reassembly is downstream
  at BgzfCompress).

Gated on a min record count so small chunks keep the cheaper serial gather.
Output is identical (same records, block boundaries, and dense ordinals),
covered by a serial-vs-parallel parity test across the count-cap, byte-cap,
single-block, and partial-tail regimes, and the existing oracle parity suite.
@nh13
nh13 deployed to github-actions September 14, 2026 22:17 — with GitHub Actions Active
@coderabbitai

coderabbitai Bot commented Sep 14, 2026 •

Copy link
Copy Markdown

Review Change StackReview Change Stack

Note

Reviews paused

Use the following commands to manage reviews:

  • @coderabbitai resume to resume automatic reviews.
  • @coderabbitai review to trigger a single review.

Use the checkboxes below for quick actions:

  • ▶️ Resume reviews
  • 🔍 Trigger review

No actionable comments were generated in the recent review. 🎉

ℹ️ Recent review info
⚙️ Run configuration

Configuration used: Path: .coderabbit.yaml

Review profile: ASSERTIVE

Plan: Essentials

Run ID: 63322bc0-b4bd-48e9-86f6-a064def75ec6

📥 Commits

Reviewing files that changed from the base of the PR and between f0455c0 and db3a77f.

📒 Files selected for processing (6)
  • crates/fgumi-pipeline-io/src/sort/merge.rs
  • crates/fgumi-pipeline-io/src/sort/protocol.rs
  • crates/fgumi-pipeline-io/src/sort/tests.rs
  • crates/fgumi-sort/src/inline.rs
  • crates/fgumi-sort/src/template_arena.rs
  • src/lib/pipeline/chains/builder.rs

Included review availability: 0 reviews are currently available. Your included PR review attempts over the past 7 days set your current allowance at 1 review per hour.


Walkthrough

Changes

Parallel in-memory sort gathering

Layer / File(s) Summary
Planning contracts and record lengths
crates/fgumi-pipeline-io/src/sort/merge.rs, crates/fgumi-pipeline-io/src/sort/protocol.rs, crates/fgumi-sort/src/inline.rs, crates/fgumi-sort/src/template_arena.rs
Output builders declare framing overhead. In-memory chunks expose record lengths for serial-equivalent block planning.
Bounded parallel fast-path execution
crates/fgumi-pipeline-io/src/sort/merge.rs
SortMerge plans bounded block windows, gathers them concurrently, preserves order, handles backpressure, and shares completion finalization.
Phase-2 thread-budget wiring
src/lib/pipeline/chains/builder.rs
Terminal and intermediate streaming sort paths pass phase-2 thread limits to in-memory fast-path gathering.
Serial and parallel parity coverage
crates/fgumi-pipeline-io/src/sort/tests.rs
Tests compare serial and parallel outputs for block framing, record batches, boundaries, ordinals, partial tails, caps, and backpressure.

Priority: ➖ Normal

Estimated code review effort: 4 (Complex) | ~45 minutes

Change: Feature

Sequence Diagram(s)

sequenceDiagram
  participant SortMerge
  participant MemoryChunk
  participant WorkerPool
  participant Output
  SortMerge->>MemoryChunk: Read record lengths
  SortMerge->>WorkerPool: Submit bounded block ranges
  WorkerPool->>Output: Produce ordered blocks or batches
  Output-->>SortMerge: Accept or reject output
  SortMerge->>Output: Resume pending output
Loading

Suggested labels: fgumi sort

Merge Risk: ⚪ Minimal · up to db3a7

The parallel fast path is mergeable based on the available evidence, with no actionable risk identified.

🚥 Pre-merge checks | ✅ 2 | ❌ 1

❌ Failed checks (1 warning)

Check name Status Explanation Resolution
Title check ⚠️ Warning The title uses a valid Conventional Commit type and clearly describes the change, but its description is not imperative; "parallel gather" is a noun phrase. Change the description to lowercase imperative wording, such as "perf: parallelize gathering for the in-memory sort fast path".
✅ Passed checks (2 passed)
Check name Status Explanation
Linked Issues check ✅ Passed Check skipped because no linked issues were found for this pull request.
Out of Scope Changes check ✅ Passed Check skipped because no linked issues were found for this pull request.
  • Fix all pre-merge checks with AI

Comment @coderabbitai help to get the list of available commands.

@nh13

nh13 commented Sep 14, 2026

Copy link
Copy Markdown
Member Author

@coderabbitai pause

@coderabbitai

coderabbitai Bot commented Sep 14, 2026

Copy link
Copy Markdown
✅ Action performed

Reviews paused.

@codecov

codecov Bot commented Sep 14, 2026 •

Copy link
Copy Markdown

Codecov Report

❌ Patch coverage is 91.13300% with 18 lines in your changes missing coverage. Please review.
✅ Project coverage is 96.05%. Comparing base (9b978a2) to head (db3a77f).
⚠️ Report is 6 commits behind head on main.

Files with missing lines Patch % Lines
crates/fgumi-pipeline-io/src/sort/merge.rs 93.25% 12 Missing ⚠️
crates/fgumi-pipeline-io/src/sort/protocol.rs 57.14% 3 Missing ⚠️
crates/fgumi-sort/src/template_arena.rs 0.00% 3 Missing ⚠️
Additional details and impacted files
@@            Coverage Diff             @@
##             main     #966      +/-   ##
==========================================
- Coverage   96.08%   96.05%   -0.03%     
==========================================
  Files         291      292       +1     
  Lines      144093   145159    +1066     
==========================================
+ Hits       138447   139435     +988     
- Misses       5646     5724      +78     

☔ View full report in Codecov by Harness.
📢 Have feedback on the report? Share it here.

🚀 New features to boost your workflow:
  • ❄️ Test Analytics: Detect flaky tests, report on failures, and find test suite problems.

@nh13

nh13 commented Sep 17, 2026

Copy link
Copy Markdown
Member Author

@coderabbitai review

@coderabbitai

coderabbitai Bot commented Sep 17, 2026 •

Copy link
Copy Markdown
✅ Action performed

Review finished.

Note: CodeRabbit is an incremental review system and does not re-review already reviewed commits. This command is applicable only when automatic reviews are paused.

@nh13
nh13 added this pull request to the merge queue Sep 17, 2026
Merged via the queue into main with commit 12db0a7 Sep 17, 2026
17 checks passed
@nh13
nh13 deleted the nh/sort-merge-parallel-gather branch September 17, 2026 16:10
@nh13 nh13 mentioned this pull request Sep 17, 2026

This branch was successfully deployed

1 active deployment
github-actions — db3a77f1 Deployed Sep 14, 2026 by nh13 via coverage #4524
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant