Skip to content

perf: Concretely typed TopK StringHeap, and StringHashTable slot reuse - #24105

Open
MassivePizza wants to merge 22 commits into
apache:mainfrom
massive-com:topk-mutable-borrow-slots
Open

MassivePizza wants to merge 22 commits into
apache:mainfrom
massive-com:topk-mutable-borrow-slots

Conversation

@MassivePizza

@MassivePizza MassivePizza commented Aug 5, 2026 •

Copy link
Copy Markdown
Contributor

Which issue does this PR close?

N/A.

Continuation of #23609.

Rationale for this change

Improve perf of TopK.

What changes are included in this PR?

Eliminate overhead of casting with dyn Any for every String in topk::StringHeap.

Rearrange HashTableItem<ID> Option-ness to avoid unnecessary unwraps on lookups in topk::TopKHashTable.
Rewrite inserts to take a borrowable type, delaying clones and allowing reuse of storage slots (mostly Strings).

Are these changes tested?

Should be covered by existing tests.

Making ID an Option<ID> inside HashTableItem should be fine, since the internal HashTable keeps track of valid items (although this could be hard to verify properly). It runs with the assumption that "if everything needed an unwrap, nothing needs an unwrap".

Are there any user-facing changes?

No.

@github-actions github-actions Bot added the physical-plan Changes to the physical-plan crate label Aug 5, 2026
@MassivePizza MassivePizza changed the title perf: Concretely typed TopK StringHeap, and improved StringHashTable slot reuse perf: Concretely typed TopK StringHeap, and StringHashTable slot reuse Aug 5, 2026
@MassivePizza
MassivePizza force-pushed the topk-mutable-borrow-slots branch from 6ed1b35 to ddd9b71 Compare August 12, 2026 20:08
@codecov-commenter

codecov-commenter commented Aug 12, 2026 •

Copy link
Copy Markdown

Codecov Report

❌ Patch coverage is 97.01493% with 12 lines in your changes missing coverage. Please review.
✅ Project coverage is 82.77%. Comparing base (791660c) to head (70dea10).

Files with missing lines Patch % Lines
.../physical-plan/src/aggregates/topk/priority_map.rs 80.00% 0 Missing and 6 partials ⚠️
...on/physical-plan/src/aggregates/topk/hash_table.rs 97.16% 1 Missing and 3 partials ⚠️
...tafusion/physical-plan/src/aggregates/topk/heap.rs 99.13% 2 Missing ⚠️
Additional details and impacted files
@@           Coverage Diff            @@
##             main   #24105    +/-   ##
========================================
  Coverage   82.76%   82.77%            
========================================
  Files        1147     1147            
  Lines      450580   450745   +165     
  Branches   450580   450745   +165     
========================================
+ Hits       372944   373108   +164     
+ Misses      54929    54927     -2     
- Partials    22707    22710     +3     

☔ 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.
  • 📦 JS Bundle Analysis: Save yourself from yourself by tracking and limiting bundle sizes in JS merges.

@MassivePizza
MassivePizza force-pushed the topk-mutable-borrow-slots branch from b6aeee8 to 4d8b2c5 Compare August 21, 2026 13:57
@MassivePizza
MassivePizza marked this pull request as ready for review August 26, 2026 15:03
@MassivePizza
MassivePizza force-pushed the topk-mutable-borrow-slots branch from 5e604aa to 5fd2a17 Compare August 27, 2026 10:30

@kosiew kosiew left a comment •

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.

@MassivePizza,

Sorry for the late review.

Thanks for working on these TopK improvements. The concrete string types and allocation reuse look like useful optimizations. I found one correctness issue around NULL group keys that I think needs to be addressed before merging. I also left a small testing suggestion for the new string heap path.

Comment thread datafusion/physical-plan/src/aggregates/topk/hash_table.rs Outdated
Comment thread datafusion/physical-plan/src/aggregates/topk/heap.rs Outdated

@kosiew kosiew left a comment

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.

@MassivePizza,

Thanks for the follow-up changes. I took another pass through the updated code. The main NULL-key correctness issue still appears to be present, and the focused TopKHeap<String> regression test is still missing.

The NULL-key issue is the one blocking approval here because it can cause a valid aggregate group to be omitted from the output. I left the details inline below.

Once that is addressed, along with the requested string heap test, this should be in much better shape.

Comment thread datafusion/physical-plan/src/aggregates/topk/hash_table.rs Outdated
Comment thread datafusion/physical-plan/src/aggregates/topk/heap.rs Outdated
@MassivePizza
MassivePizza requested a review from kosiew October 7, 2026 15:55
@kosiew

kosiew commented Oct 9, 2026

Copy link
Copy Markdown
Contributor

run benchmark topk_aggregate

@adriangbot

Copy link
Copy Markdown

Benchmark for this request failed before finishing (Kubernetes reason: BackoffLimitExceeded).

Benchmarks requested: topk_aggregate

Runner log (last 40 lines)
2026-10-09T04:26:14.425995Z  INFO runner starting benchmark runner bench_type=Datafusion, pr_url=https://github.com/apache/datafusion/pull/24105, benchmarks=topk_aggregate
2026-10-09T04:26:14.484427Z  INFO benchmark_controller::runner::bench_datafusion === Cloning PR branch ===
2026-10-09T04:26:14.484447Z  INFO benchmark_controller::runner::shell running command cmd=git, args=["clone", "--depth=200", "https://github.com/apache/datafusion.git", "/workspace/datafusion-branch"], cwd="/"
2026-10-09T04:26:19.494030Z  INFO benchmark_controller::runner::shell running command cmd=git, args=["fetch", "origin", "refs/pull/24105/head:topk-mutable-borrow-slots", "main"], cwd="/workspace/datafusion-branch"
2026-10-09T04:26:39.499400Z  INFO benchmark_controller::runner::shell running command cmd=git, args=["checkout", "topk-mutable-borrow-slots"], cwd="/workspace/datafusion-branch"
2026-10-09T04:26:44.502224Z  INFO benchmark_controller::runner::shell running command cmd=git, args=["merge-base", "HEAD", "origin/main"], cwd="/workspace/datafusion-branch"
2026-10-09T04:26:49.504165Z ERROR runner benchmark failed error=git merge-base: git merge-base HEAD origin/main exited with code 1
stdout:

stderr:
Kubernetes message
Job has reached the specified backoff limit

File an issue against this benchmark runner

@kosiew

kosiew commented Oct 9, 2026

Copy link
Copy Markdown
Contributor

run benchmark topk_aggregate

@adriangbot

Copy link
Copy Markdown

🤖 Benchmark running (GKE) | trigger
Instance: c4a-highmem-16 (12 vCPU / 65 GiB) | Linux bench-c6074389145-3227-b9ggd 6.12.94+ #1 SMP Fri Aug 21 08:00:16 UTC 2026 aarch64 GNU/Linux

CPU Details (lscpu)
Architecture:                            aarch64
CPU op-mode(s):                          64-bit
Byte Order:                              Little Endian
CPU(s):                                  16
On-line CPU(s) list:                     0-15
Vendor ID:                               ARM
Model name:                              Neoverse-V2
Model:                                   1
Thread(s) per core:                      1
Core(s) per cluster:                     16
Socket(s):                               -
Cluster(s):                              1
Stepping:                                r0p1
BogoMIPS:                                2000.00
Flags:                                   fp asimd evtstrm aes pmull sha1 sha2 crc32 atomics fphp asimdhp cpuid asimdrdm jscvt fcma lrcpc dcpop sha3 sm3 sm4 asimddp sha512 sve asimdfhm dit uscat ilrcpc flagm sb paca pacg dcpodp sve2 sveaes svepmull svebitperm svesha3 svesm4 flagm2 frint svei8mm svebf16 i8mm bf16 dgh rng bti
L1d cache:                               1 MiB (16 instances)
L1i cache:                               1 MiB (16 instances)
L2 cache:                                32 MiB (16 instances)
L3 cache:                                80 MiB (1 instance)
NUMA node(s):                            1
NUMA node0 CPU(s):                       0-15
Vulnerability Gather data sampling:      Not affected
Vulnerability Indirect target selection: Not affected
Vulnerability Itlb multihit:             Not affected
Vulnerability L1tf:                      Not affected
Vulnerability Mds:                       Not affected
Vulnerability Meltdown:                  Not affected
Vulnerability Mmio stale data:           Not affected
Vulnerability Reg file data sampling:    Not affected
Vulnerability Retbleed:                  Not affected
Vulnerability Spec rstack overflow:      Not affected
Vulnerability Spec store bypass:         Mitigation; Speculative Store Bypass disabled via prctl
Vulnerability Spectre v1:                Mitigation; __user pointer sanitization
Vulnerability Spectre v2:                Mitigation; CSV2, BHB
Vulnerability Srbds:                     Not affected
Vulnerability Tsa:                       Not affected
Vulnerability Tsx async abort:           Not affected
Vulnerability Vmscape:                   Not affected

Comparing topk-mutable-borrow-slots (70dea10) to 791660c (merge-base) diff

Run configuration
run benchmark topk_aggregate

Results will be posted here when complete


File an issue against this benchmark runner

@adriangbot

Copy link
Copy Markdown

🤖 Benchmark completed (GKE) | trigger

Instance: c4a-highmem-16 (12 vCPU / 65 GiB)

Comparing topk-mutable-borrow-slots (70dea10) to 791660c (merge-base) diff

Run configuration
run benchmark topk_aggregate
CPU Details (lscpu)
Architecture:                            aarch64
CPU op-mode(s):                          64-bit
Byte Order:                              Little Endian
CPU(s):                                  16
On-line CPU(s) list:                     0-15
Vendor ID:                               ARM
Model name:                              Neoverse-V2
Model:                                   1
Thread(s) per core:                      1
Core(s) per cluster:                     16
Socket(s):                               -
Cluster(s):                              1
Stepping:                                r0p1
BogoMIPS:                                2000.00
Flags:                                   fp asimd evtstrm aes pmull sha1 sha2 crc32 atomics fphp asimdhp cpuid asimdrdm jscvt fcma lrcpc dcpop sha3 sm3 sm4 asimddp sha512 sve asimdfhm dit uscat ilrcpc flagm sb paca pacg dcpodp sve2 sveaes svepmull svebitperm svesha3 svesm4 flagm2 frint svei8mm svebf16 i8mm bf16 dgh rng bti
L1d cache:                               1 MiB (16 instances)
L1i cache:                               1 MiB (16 instances)
L2 cache:                                32 MiB (16 instances)
L3 cache:                                80 MiB (1 instance)
NUMA node(s):                            1
NUMA node0 CPU(s):                       0-15
Vulnerability Gather data sampling:      Not affected
Vulnerability Indirect target selection: Not affected
Vulnerability Itlb multihit:             Not affected
Vulnerability L1tf:                      Not affected
Vulnerability Mds:                       Not affected
Vulnerability Meltdown:                  Not affected
Vulnerability Mmio stale data:           Not affected
Vulnerability Reg file data sampling:    Not affected
Vulnerability Retbleed:                  Not affected
Vulnerability Spec rstack overflow:      Not affected
Vulnerability Spec store bypass:         Mitigation; Speculative Store Bypass disabled via prctl
Vulnerability Spectre v1:                Mitigation; __user pointer sanitization
Vulnerability Spectre v2:                Mitigation; CSV2, BHB
Vulnerability Srbds:                     Not affected
Vulnerability Tsa:                       Not affected
Vulnerability Tsx async abort:           Not affected
Vulnerability Vmscape:                   Not affected
Details

group                                                             HEAD                                   topk-mutable-borrow-slots
-----                                                             ----                                   -------------------------
aggregate 10000000 time-series rows                               1.00    109.9±4.74ms        ? ?/sec    1.01    110.7±5.75ms        ? ?/sec
aggregate 10000000 worst-case rows                                1.00    110.5±5.13ms        ? ?/sec    1.00    110.6±5.91ms        ? ?/sec
distinct 10000000 rows asc [TopK]                                 1.07      4.2±0.08ms        ? ?/sec    1.00      3.9±0.08ms        ? ?/sec
distinct 10000000 rows asc [no TopK]                              1.00     52.2±2.65ms        ? ?/sec    1.00     52.4±2.58ms        ? ?/sec
distinct 10000000 rows desc [TopK]                                1.09      4.3±0.10ms        ? ?/sec    1.00      3.9±0.06ms        ? ?/sec
distinct 10000000 rows desc [no TopK]                             1.00     53.2±2.76ms        ? ?/sec    1.00     53.2±2.33ms        ? ?/sec
string aggregate 10000000 time-series rows [Utf8View]             1.07     51.0±1.74ms        ? ?/sec    1.00     47.5±1.09ms        ? ?/sec
string aggregate 10000000 time-series rows [Utf8]                 1.07     49.1±1.38ms        ? ?/sec    1.00     45.7±0.99ms        ? ?/sec
string aggregate 10000000 worst-case rows [Utf8View]              1.13   351.2±14.62ms        ? ?/sec    1.00   310.6±11.35ms        ? ?/sec
string aggregate 10000000 worst-case rows [Utf8]                  1.11   335.4±12.33ms        ? ?/sec    1.00   303.4±11.05ms        ? ?/sec
top k=10 aggregate 10000000 time-series rows                      1.15      9.3±0.25ms        ? ?/sec    1.00      8.1±0.25ms        ? ?/sec
top k=10 aggregate 10000000 time-series rows [Utf8View]           1.11     10.6±0.38ms        ? ?/sec    1.00      9.5±0.28ms        ? ?/sec
top k=10 aggregate 10000000 worst-case rows                       1.25     14.8±0.82ms        ? ?/sec    1.00     11.9±0.46ms        ? ?/sec
top k=10 aggregate 10000000 worst-case rows [Utf8View]            1.18     16.0±0.75ms        ? ?/sec    1.00     13.5±0.48ms        ? ?/sec
top k=10 string aggregate 10000000 time-series rows [Utf8View]    1.30     11.5±0.33ms        ? ?/sec    1.00      8.9±0.26ms        ? ?/sec
top k=10 string aggregate 10000000 time-series rows [Utf8]        1.39      9.8±0.23ms        ? ?/sec    1.00      7.1±0.17ms        ? ?/sec
top k=10 string aggregate 10000000 worst-case rows [Utf8View]     1.29     11.5±0.26ms        ? ?/sec    1.00      8.9±0.26ms        ? ?/sec
top k=10 string aggregate 10000000 worst-case rows [Utf8]         1.38      9.8±0.15ms        ? ?/sec    1.00      7.1±0.19ms        ? ?/sec

Resource Usage

topk_aggregate — base (merge-base)

Metric Value
Wall time 720.2s
Peak memory 2.1 GiB
Avg memory 505.1 MiB
CPU user 2103.4s
CPU sys 136.2s
Peak spill 0 B

topk_aggregate — branch

Metric Value
Wall time 695.2s
Peak memory 2.1 GiB
Avg memory 486.9 MiB
CPU user 1928.3s
CPU sys 137.8s
Peak spill 0 B

File an issue against this benchmark runner

@kosiew kosiew left a comment

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.

@MassivePizza,

Thanks for addressing the previous review comments. I've reviewed the follow-up changes and confirmed that both issues are resolved.

The separate NULL and vacant-slot sentinels correctly distinguish NULL grouping keys from freed entries, and the new regression tests cover NULL-key emission and the string heap's insertion, replacement, and drain paths.

The NULL-index collection also preserves key/value alignment during emission without requiring an intermediate vector.

I have no further concerns. LGTM!

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

Labels

physical-plan Changes to the physical-plan crate

Projects

None yet

Development

Successfully merging this pull request may close these issues.

4 participants