Skip to content

feat: reduce open file handles during IVF training - #6169

Merged
westonpace merged 2 commits into
lance-format:mainfrom
westonpace:feat/two-file-shuffler
Mar 16, 2026
Merged

westonpace merged 2 commits into
lance-format:mainfrom
westonpace:feat/two-file-shuffler

Conversation

@westonpace

Copy link
Copy Markdown
Member

The previous shuffler used one open file per partition. At large scales this meant tedious re-adjusting of OS limits. The new shuffler uses two open files. We potentially introduce a bit more random access in the later read phase but the overall performance hasn't changed significantly.

@github-actions github-actions Bot added the enhancement New feature or request label Mar 11, 2026
@github-actions

Copy link
Copy Markdown
Contributor

PR Review

Nice improvement — reducing open file descriptors from O(partitions) to O(1) is a meaningful usability win for large-scale IVF training. The two-file design with sorted data + cumulative offsets is clean and well-tested. A few items to consider:

P1: Potential u32 overflow in partition_ranges

In partition_ranges (shuffler.rs), the offset position is computed as:

let end_pos = (batch_idx as usize * self.num_partitions + partition_id) as u32;

This silently truncates to u32 if num_batches * num_partitions > u32::MAX. With the default 128MB batch size, processing ~10TB of data across 100K+ partitions would overflow. Since this PR specifically targets users with many partitions, consider using u64 positions (if ReadBatchParams supports it) or at minimum adding a checked cast / early validation that the offsets file won't exceed u32 rows.

Minor: precomputed_shuffle_buffers dead code

TwoFileShuffler declares the precomputed_shuffle_buffers field and with_precomputed_shuffle_buffers setter but never uses them (#[allow(dead_code)]). If this feature isn't planned for the new shuffler, removing the field/method avoids confusion. If it is planned, a tracking comment or TODO would help.

Tests

Good coverage of the core paths: round-trip, empty partitions, loss tracking, multi-batch spilling. The with_batch_size_bytes(16) test for forcing multiple batches is a nice touch.

🤖 Generated with Claude Code

@codecov

codecov Bot commented Mar 11, 2026 •

Copy link
Copy Markdown

Codecov Report

❌ Patch coverage is 91.79811% with 26 lines in your changes missing coverage. Please review.

Files with missing lines Patch % Lines
rust/lance-index/src/vector/v3/shuffler.rs 91.77% 14 Missing and 12 partials ⚠️

📢 Thoughts on this report? Let us know!

let offsets = offsets_stream.try_collect::<Vec<_>>().await?;
let offsets = if offsets.len() == 0 {
// We should not hit this path if there is no batches
unreachable!()

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

should it panic with a message?

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

unreachable panics with the message "I didn't expect to get here" and the line number which is pretty much what I'd use if I did an explicit panic.


let mut ranges = Vec::with_capacity(self.num_batches as usize);
for batch_idx in 0..self.num_batches {
if batch_idx == 0 && is_uneven {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

is the is_uneven stuff equivalent to checking if partition_id = 0 and batch_idx = 0?

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yes, that's maybe simpler. Updated.

@wkalt

wkalt commented Mar 11, 2026 •

Copy link
Copy Markdown
Contributor
wpace_progress_full

here is the result on my test of 500M rows for posterity.

The previous behavior looked like this:
progress

@westonpace

Copy link
Copy Markdown
Member Author

P1: Potential u32 overflow in partition_ranges

This is a legitimate concern. Unfortunately, we don't actually support take with u64 offsets. This is because we're reusing some old paths from the v0.1 days where we used u32 offsets. It's all fixable but more follow-up territory.

By my math I think, even with 1536-dimension vectors, we are ok until we hit trillions of rows.

1T rows => sqrt(1T) * 96 * 1T bytes of PQ data => 96EB of data which, split into 128MB chunks, would give us ~1B chunks.

Either way, I turned it into a try_into so we get a nice error and can address in the future.

@Xuanwo Xuanwo left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thank you for this PR! Only two small questions.

///
/// This default is likely to be fine for most use cases.
fn shuffle_batch_bytes() -> usize {
std::env::var("LANCE_SHUFFLE_BATCH_BYTES")

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I think it's a good idea to add some protection here. This would prevent users from setting LANCE_SHUFFLE_BATCH_BYTES to 0, which could lead to using batch_size_bytes = 0 and producing incorrect results.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Good idea, now it will log a warning and use the default.

let mut partition_counts = vec![0u64; np];
for i in 0..part_ids.len() {
let pid = part_ids.value(i) as usize;
if pid < np {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Is it possible for pid >= np? It looks like we will just write those data without offsets. global_row_count will always include them.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

It shouldn't be possible (np is num_partitions) but I now log a warning instead of silently ignoring.

@westonpace
westonpace force-pushed the feat/two-file-shuffler branch from a692013 to 203b465 Compare March 13, 2026 13:38
…num_partitions file handles

Actually write the offsets to a file and don't accumulate

Change progress reporting to report number of rows shuffled and not number of batches processed

Address review suggestions

More PR suggestions

Remove dead code

Address PR review
@westonpace
westonpace force-pushed the feat/two-file-shuffler branch from 203b465 to 08715e5 Compare March 16, 2026 14:05
@westonpace

Copy link
Copy Markdown
Member Author

CI failure seems unrelated. Will merge if remaining CI job passes.

@westonpace
westonpace merged commit bfe7f56 into lance-format:main Mar 16, 2026
27 of 28 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

enhancement New feature or request

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants