Skip to content

Don't split files whose index selection is a handful of rows - #10360

Open
joseph-isaacs wants to merge 2 commits into
developfrom
ji/sparse-selection-partitions
Open

joseph-isaacs wants to merge 2 commits into
developfrom
ji/sparse-selection-partitions

Conversation

@joseph-isaacs

@joseph-isaacs joseph-isaacs commented Oct 7, 2026 •

Copy link
Copy Markdown
Contributor

Summary

When an external index has already picked the rows of each file, DataFusion still splits every file into byte ranges, one per partition. Each partition opens the file and plans a scan, then finds that its range holds none of the selected rows. With one 14M-row file and 4 partitions, a 6-row lookup paid for four opens and used one.

Changes

  • VortexSource::repartitioned keeps files whole when every file in the scan carries an inclusive selection (IncludeByIndex or IncludeRoaring) of at most SPARSE_SELECTION_MAX_ROWS rows (8,192, one batch). Files are already spread across file groups, so nothing else changes. Larger selections still split, because scattered rows benefit from parallel partitions: a 100K threshold made a 58K-row selection slower (2.9 → 3.8 ms).
  • Selection::row_mask for IncludeByIndex sets the mask bits in place instead of first collecting relative indices into a Vec<usize>. A split whose rows are all selected becomes an all-true mask.

New test: an rstest covering a file with no access plan, a sparse selection, a selection of exactly one batch, and a dense selection.

End to end on a 14M-row file (DataFusion, 4 partitions, index bitmap → access plan → count(*)), CPU per query, two runs each. Before and after differ only by this change:

Query Rows selected CPU before CPU after
trace_id lookup 6 1.71 / 1.74 ms 1.59 / 1.54 ms
rare token lookup 1 1.67 / 1.72 ms 1.58 / 1.58 ms
common token 840K 9.4 / 10.0 ms 8.6 / 7.9 ms

When an external index pins down the rows of every file, DataFusion still
splits each file into byte ranges, one per partition. Every partition then
opens the file, plans the scan and finds that its range holds none of the
selected rows. `VortexSource::repartitioned` now keeps files whole when
every file carries an inclusive selection of at most one batch of rows;
larger selections still split for parallelism.

`Selection::row_mask` for `IncludeByIndex` also stops collecting relative
indices into a `Vec<usize>` before building the mask, setting the bits in
place instead.

End to end on a 14M-row file (index bitmap in, SQL out, 4 partitions):
6-row count CPU 2.2 ms -> 1.55 ms, 840K-row count 11.5 ms -> 8.2 ms.

Signed-off-by: Joe Isaacs <joe.isaacs@live.co.uk>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01JRTKSkM9xUFT8WHRyS97aR
@codspeed

codspeed Bot commented Oct 7, 2026 •

Copy link
Copy Markdown

Merging this PR will regress 1 benchmark

⚠️ Unknown Walltime execution environment detected

Using the Walltime instrument on standard Hosted Runners will lead to inconsistent data.

For the most accurate results, we recommend using CodSpeed Macro Runners: bare-metal machines fine-tuned for performance measurement consistency.

⚠️ Different runtime environments detected

Some benchmarks with significant performance changes were compared across different runtime environments,
which may affect the accuracy of the results.

Open the report in CodSpeed to investigate

⚡ 2 improved benchmarks
❌ 1 regressed benchmark
✅ 2120 untouched benchmarks
⏩ 518 skipped benchmarks1

Warning

Please fix the performance issues or acknowledge them on CodSpeed.

Performance Changes

Mode Benchmark BASE HEAD Efficiency
❌ WallTime bitpack_blocked_compress_avx2 6.7 µs 7.6 µs -11.85%
⚡ WallTime bitpack_blocked_compress_avx512 5.5 µs 4.7 µs +17.24%
⚡ Simulation column_x_column_points 236.6 µs 213.9 µs +10.62%

Tip

Investigate this regression by commenting @codspeedbot fix this regression on this PR, or directly use the CodSpeed MCP with your agent.


Comparing ji/sparse-selection-partitions (c339cc2) with develop (438600d)2

Open in CodSpeed

Footnotes

  1. 518 benchmarks were skipped, so the baseline results were used instead. If they were deleted from the codebase, click here and archive them to remove them from the performance reports. ↩

  2. No successful run was found on develop (09df3e8) during the generation of this report, so 438600d was used instead as the comparison base. There might be some changes unrelated to this pull request in this report. ↩

@joseph-isaacs joseph-isaacs added the changelog/performance A performance improvement label Oct 7, 2026
Move the scan over every file's access plan into
`every_file_has_sparse_selection`, and say what returning `None` means to
DataFusion and that the fallback mirrors DataFusion's default
repartitioning, which an override cannot call. No behaviour change.

Signed-off-by: Joe Isaacs <joe.isaacs@live.co.uk>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01JRTKSkM9xUFT8WHRyS97aR

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

changelog/performance A performance improvement

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants