Re: [PR] refactor: Update SortMergeJoin to use async spill abstractions [datafusion]
pantShrey commented on PR #22230: URL: https://github.com/apache/datafusion/pull/22230#issuecomment-4718662614 @kumarUjjawal Thank you so much for taking the time to review and merge this PR! @mbutrovich, thank you as well for your review and feedback. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] - To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
Re: [PR] refactor: Update SortMergeJoin to use async spill abstractions [datafusion]
kumarUjjawal merged PR #22230: URL: https://github.com/apache/datafusion/pull/22230 -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] - To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
Re: [PR] refactor: Update SortMergeJoin to use async spill abstractions [datafusion]
kumarUjjawal commented on PR #22230: URL: https://github.com/apache/datafusion/pull/22230#issuecomment-4718432306 Thank you all -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] - To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
Re: [PR] refactor: Update SortMergeJoin to use async spill abstractions [datafusion]
adriangbot commented on PR #22230: URL: https://github.com/apache/datafusion/pull/22230#issuecomment-4710239478 π€ Benchmark completed (GKE) | [trigger](https://github.com/apache/datafusion/pull/22230#issuecomment-4710067506) **Instance:** `c4a-highmem-16` (12 vCPU / 65 GiB) 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 ``` Comparing HEAD and smj-async-spill Benchmark tpch_sf10.json βββββ³β³β³ββββ β Query β HEAD β smj-async-spill βChange β β‘ββββββββββ© β QQuery 1 β 314.02 / 317.61 Β±3.33 / 323.43 ms β 310.14 / 312.17 Β±2.29 / 316.39 ms β no change β β QQuery 2 β 101.27 / 105.78 Β±3.83 / 112.57 ms β 100.14 / 105.01 Β±3.59 / 109.85 ms β no change β β QQuery 3 β 236.67 / 240.51 Β±3.73 / 247.34 ms β 233.85 / 240.00 Β±3.57 / 244.03 ms β no change β β QQuery 4 β 113.82 / 116.37 Β±2.68 / 121.00 ms β 113.48 / 116.61 Β±2.97 / 122.24 ms β no change β β QQuery 5 β 347.18 / 358.05 Β±9.35 / 371.55 ms β 349.35 / 355.67 Β±5.67 / 364.77 ms β no change β β QQuery 6 β 124.27 / 128.85 Β±3.61 / 132.57 ms β 124.46 / 129.85 Β±4.16 / 134.31 ms β no change β β QQuery 7 β 459.21 / 467.46 Β±7.69 / 480.72 ms β 460.80 / 467.42 Β±5.48 / 473.94 ms β no change β β QQuery 8 β 378.69 / 385.84 Β±4.53 / 391.29 ms β 390.61 / 394.59 Β±2.62 / 398.69 ms β no change β β QQuery 9 β 532.28 / 548.90 Β±13.38 / 563.28 ms β 533.96 / 549.20 Β±12.22 / 565.85 ms β no change β β QQuery 10 β 304.28 / 319.37 Β±12.76 / 337.55 ms β 305.59 / 319.81 Β±18.17 / 354.08 ms β no change β β QQuery 11 β84.22 / 93.71 Β±8.42 / 108.28 ms β 87.29 / 95.28 Β±10.16 / 114.92 ms β no change β β QQuery 12 β 177.70 / 183.80 Β±5.70 / 193.73 ms β 175.60 / 184.84 Β±13.24 / 211.11 ms β no change β β QQuery 13 β 287.19 / 301.67 Β±9.71 / 316.71 ms β 283.74 / 294.79 Β±6.91 / 302.87 ms β no change β β QQuery 14 β 173.58 / 177.24 Β±4.39 / 185.19 ms β 173.15 / 176.27 Β±2.56 / 180.45 ms β no change β β QQuery 15 β 305.13 / 310.44 Β±5.60 / 321.29 ms β 314.36 / 320.79 Β±6.42 / 329.86 ms β no change β β QQuery 16 β 68.75 / 71.03 Β±2.10 / 74.61 ms β 67.69 / 69.28 Β±2.03 / 73.28 ms β no change β β QQuery 17 β 635.95 / 675.8
Re: [PR] refactor: Update SortMergeJoin to use async spill abstractions [datafusion]
adriangbot commented on PR #22230: URL: https://github.com/apache/datafusion/pull/22230#issuecomment-4710099940 π€ Benchmark running (GKE) | [trigger](https://github.com/apache/datafusion/pull/22230#issuecomment-4710067506) **Instance:** `c4a-highmem-16` (12 vCPU / 65 GiB) | `Linux bench-c4710067506-571-8rhh5 6.12.68+ #1 SMP Sat May 2 07:49:07 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 smj-async-spill (81198d7ea38dacc11e950b29bee57b8c3b9d7a89) to 666f862 (merge-base) [diff](https://github.com/apache/datafusion/compare/666f86210e0e03f1207cf94a8718d5b0fdb2270d..81198d7ea38dacc11e950b29bee57b8c3b9d7a89) using: tpch10 Results will be posted here when complete --- [File an issue](https://github.com/adriangb/datafusion-benchmarking/issues) against this benchmark runner -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] - To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
Re: [PR] refactor: Update SortMergeJoin to use async spill abstractions [datafusion]
kumarUjjawal commented on PR #22230: URL: https://github.com/apache/datafusion/pull/22230#issuecomment-4710067506 run benchmark tpch10 ``` env: PREFER_HASH_JOIN: false ``` -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] - To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
Re: [PR] refactor: Update SortMergeJoin to use async spill abstractions [datafusion]
github-actions[bot] commented on PR #22230: URL: https://github.com/apache/datafusion/pull/22230#issuecomment-4708756271 Thank you for opening this pull request! Reviewer note: [cargo-semver-checks](https://github.com/obi1kenobi/cargo-semver-checks) reported the current version number is not SemVer-compatible with the changes in this pull request (compared against the base branch). Details ``` Cloning apache/main Building datafusion-physical-plan v54.0.0 (current) Built [ 36.958s] (current) Parsing datafusion-physical-plan v54.0.0 (current) Parsed [ 0.128s] (current) Building datafusion-physical-plan v54.0.0 (baseline) Built [ 36.115s] (baseline) Parsing datafusion-physical-plan v54.0.0 (baseline) Parsed [ 0.133s] (baseline) Checking datafusion-physical-plan v54.0.0 -> v54.0.0 (no change; assume patch) Checked [ 0.552s] 223 checks: 222 pass, 1 fail, 0 warn, 30 skip --- failure method_parameter_count_changed: pub method parameter count changed --- Description: A publicly-visible method now takes a different number of parameters, not counting the receiver (self) parameter. ref: https://doc.rust-lang.org/cargo/reference/semver.html#fn-change-arity impl: https://github.com/obi1kenobi/cargo-semver-checks/tree/v0.48.0/src/lints/method_parameter_count_changed.ron Failed in: datafusion_physical_plan::aggregates::AggregateExec::compute_properties takes 7 parameters in /home/runner/work/datafusion/datafusion/target/semver-checks/git-apache_main/9324c2f1ef38d62cf90c6e45f4e1d7a5a00ee27f/datafusion/physical-plan/src/aggregates/mod.rs:1116, but now takes 6 parameters in /home/runner/work/datafusion/datafusion/datafusion/physical-plan/src/aggregates/mod.rs:1114 Summary semver requires new major version: 1 major and 0 minor checks failed Finished [ 75.351s] datafusion-physical-plan ``` -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] - To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
Re: [PR] refactor: Update SortMergeJoin to use async spill abstractions [datafusion]
pantShrey commented on PR #22230: URL: https://github.com/apache/datafusion/pull/22230#issuecomment-4671316486 @mbutrovich Thank you so much for the review! I have opened #22879 to track the Poll::Pending spill stream test. Will add it as a fast-follow once #21882 lands. Also addressed the `spill_stream_has_data` reset and added a peak memory bound check to the spilling SMJ test. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] - To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
Re: [PR] refactor: Update SortMergeJoin to use async spill abstractions [datafusion]
mbutrovich commented on code in PR #22230:
URL: https://github.com/apache/datafusion/pull/22230#discussion_r3383772961
##
datafusion/physical-plan/src/joins/sort_merge_join/materializing_stream.rs:
##
@@ -917,6 +943,81 @@ impl MaterializingSortMergeJoinStream {
Poll::Pending
}
+/// Identifies which buffered batches are needed for the upcoming freeze
operation
+fn get_required_batch_indices(&self, buffered_freeze_count: usize) ->
Vec {
+let mut needed = vec![];
+
+// We need all batches that matched with streamed rows
+for chunk in &self.streamed_batch.output_indices {
+if let Some(idx) = chunk.buffered_batch_idx {
+needed.push(idx);
+}
+}
+
+// Full Joins need to emit null-joined rows, so we need batches up to
freeze_count
+if self.join_type == JoinType::Full {
+needed.extend(0..buffered_freeze_count);
+}
+
+needed.sort_unstable();
+needed.dedup();
+needed
+}
+
+/// Asynchronously reads spilled batches back into memory.
+/// Only processes the required indices to avoid OOMs.
+fn poll_spilled_batches(
+&mut self,
+cx: &mut Context<'_>,
+required_indices: &[usize],
+) -> Poll> {
+for &idx in required_indices {
+// Guard against indices that might be out of bounds if the queue
was cleared
+if idx >= self.buffered_data.batches.len() {
+continue;
+}
+
+let bb = &mut self.buffered_data.batches[idx];
+
+if let BufferedBatchState::Spilled(spill_file) = &bb.batch {
+if self.spill_stream.is_none() {
+let stream = self
+.spill_manager
+.read_spill_as_stream(spill_file.clone(), None)?;
+self.spill_stream = Some(stream);
+}
+
+match
ready!(self.spill_stream.as_mut().unwrap().poll_next_unpin(cx)) {
Review Comment:
Could you open a tracking issue for this one? I don't want it to get lost
and it's the only way to really exercise the most important logic in this PR.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
-
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]
Re: [PR] refactor: Update SortMergeJoin to use async spill abstractions [datafusion]
adriangbot commented on PR #22230: URL: https://github.com/apache/datafusion/pull/22230#issuecomment-4663753385 π€ Benchmark completed (GKE) | [trigger](https://github.com/apache/datafusion/pull/22230#issuecomment-4663594101) **Instance:** `c4a-highmem-16` (12 vCPU / 65 GiB) 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 ``` Comparing HEAD and smj-async-spill Benchmark tpch_sf10.json βββββ³β³β³ββββ β Query β HEAD β smj-async-spill βChange β β‘ββββββββββ© β QQuery 1 β 317.61 / 320.09 Β±2.70 / 325.34 ms β 315.93 / 316.92 Β±1.65 / 320.21 ms β no change β β QQuery 2 β 100.91 / 104.43 Β±2.17 / 107.70 ms β 101.48 / 105.68 Β±2.58 / 109.22 ms β no change β β QQuery 3 β 239.97 / 242.81 Β±2.02 / 245.57 ms β 250.27 / 254.46 Β±2.62 / 257.45 ms β no change β β QQuery 4 β 118.41 / 119.66 Β±1.28 / 121.98 ms β 119.40 / 122.40 Β±3.11 / 127.21 ms β no change β β QQuery 5 β 371.98 / 375.82 Β±1.96 / 377.37 ms β 365.48 / 374.98 Β±6.17 / 383.10 ms β no change β β QQuery 6 β 127.11 / 129.31 Β±2.75 / 134.75 ms β 127.44 / 130.02 Β±2.35 / 134.45 ms β no change β β QQuery 7 β 484.23 / 493.57 Β±8.51 / 507.27 ms β 467.32 / 488.61 Β±12.30 / 501.72 ms β no change β β QQuery 8 β 398.15 / 405.69 Β±6.09 / 414.09 ms β 392.24 / 397.39 Β±3.06 / 400.25 ms β no change β β QQuery 9 β 559.20 / 578.21 Β±9.68 / 585.20 ms β 558.39 / 569.94 Β±11.77 / 592.11 ms β no change β β QQuery 10 β 306.83 / 321.64 Β±9.90 / 333.59 ms β 316.30 / 319.77 Β±3.57 / 326.22 ms β no change β β QQuery 11 β 85.73 / 92.28 Β±4.41 / 98.17 ms β88.80 / 92.62 Β±5.41 / 103.34 ms β no change β β QQuery 12 β 178.87 / 188.86 Β±12.68 / 212.68 ms β 184.19 / 188.72 Β±5.20 / 198.44 ms β no change β β QQuery 13 β 290.66 / 310.94 Β±13.69 / 330.06 ms β 292.23 / 309.03 Β±11.06 / 322.10 ms β no change β β QQuery 14 β 178.36 / 183.83 Β±5.21 / 193.06 ms β 178.75 / 183.58 Β±4.34 / 189.86 ms β no change β β QQuery 15 β 316.70 / 322.21 Β±4.62 / 329.28 ms β 317.78 / 321.40 Β±3.99 / 328.48 ms β no change β β QQuery 16 β 66.62 / 71.70 Β±5.84 / 82.86 ms β 66.46 / 70.10 Β±2.48 / 72.78 ms β no change β β QQuery 17 β 640.83 / 665.1
Re: [PR] refactor: Update SortMergeJoin to use async spill abstractions [datafusion]
adriangbot commented on PR #22230: URL: https://github.com/apache/datafusion/pull/22230#issuecomment-4663617675 π€ Benchmark running (GKE) | [trigger](https://github.com/apache/datafusion/pull/22230#issuecomment-4663594101) **Instance:** `c4a-highmem-16` (12 vCPU / 65 GiB) | `Linux bench-c4663594101-524-mv7bz 6.12.68+ #1 SMP Sat May 2 07:49:07 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 smj-async-spill (b7a54b192146d1fcfa00ffdedf304faa38914002) to 84bc876 (merge-base) [diff](https://github.com/apache/datafusion/compare/84bc8761ac3a126e41658b6cd0ec6bd8cc34cda8..b7a54b192146d1fcfa00ffdedf304faa38914002) using: tpch10 Results will be posted here when complete --- [File an issue](https://github.com/adriangb/datafusion-benchmarking/issues) against this benchmark runner -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] - To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
Re: [PR] refactor: Update SortMergeJoin to use async spill abstractions [datafusion]
mbutrovich commented on PR #22230: URL: https://github.com/apache/datafusion/pull/22230#issuecomment-4663594101 run benchmark tpch10 ``` env: PREFER_HASH_JOIN: false ``` -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] - To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
Re: [PR] refactor: Update SortMergeJoin to use async spill abstractions [datafusion]
mbutrovich commented on PR #22230: URL: https://github.com/apache/datafusion/pull/22230#issuecomment-4632288588 > @mbutrovich or @2010YOUY01 wonder if you have time to review this PR (at least give us some hints about how to run the relevant performance tests for sort merge join) to make sure we aren't messing up erformance > > The larger context is to try and make it easier to extend spilling (so not always to local files) We can request tpch/tpcds with prefer hash join set to false so it runs in smj. I'll do that shortly and take a review pass on this today. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] - To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
Re: [PR] refactor: Update SortMergeJoin to use async spill abstractions [datafusion]
adriangbot commented on PR #22230: URL: https://github.com/apache/datafusion/pull/22230#issuecomment-4619899638 π€ Criterion benchmark running (GKE) | [trigger](https://github.com/apache/datafusion/pull/22230#issuecomment-4619877677) **Instance:** `c4a-highmem-16` (12 vCPU / 65 GiB) | `Linux bench-c4619877677-430-tdj8l 6.12.68+ #1 SMP Wed Apr 1 02:23:28 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 smj-async-spill (8427ca087d7f0275d4ba783f52a0d59cb4ef53b3) to d2d0357 (merge-base) [diff](https://github.com/apache/datafusion/compare/d2d0357ecce85506c2aa55765b159cff27cce10e..8427ca087d7f0275d4ba783f52a0d59cb4ef53b3) BENCH_NAME=sort_merge_join BENCH_COMMAND=cargo bench --features=parquet --bench sort_merge_join BENCH_FILTER= Results will be posted here when complete --- [File an issue](https://github.com/adriangb/datafusion-benchmarking/issues) against this benchmark runner -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] - To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
Re: [PR] refactor: Update SortMergeJoin to use async spill abstractions [datafusion]
adriangbot commented on PR #22230: URL: https://github.com/apache/datafusion/pull/22230#issuecomment-4619900222 Benchmark for [this request](https://github.com/apache/datafusion/pull/22230#issuecomment-4619877677) failed. Last 20 lines of output: Click to expand ``` Cloning into '/workspace/datafusion-branch'... From https://github.com/apache/datafusion * [new ref] refs/pull/22230/head -> smj-async-spill * branchmain -> FETCH_HEAD Switched to branch 'smj-async-spill' d2d0357ecce85506c2aa55765b159cff27cce10e Cloning into '/workspace/datafusion-base'... HEAD is now at d2d0357 Port LikeExpr to use try_to_proto / try_from_proto (#22471) rustc 1.95.0 (59807616e 2026-04-14) 8427ca087d7f0275d4ba783f52a0d59cb4ef53b3 d2d0357ecce85506c2aa55765b159cff27cce10e Blocking waiting for file lock on package cache Blocking waiting for file lock on package cache Blocking waiting for file lock on package cache error: target `sort_merge_join` in package `datafusion-physical-plan` requires the features: `test_utils` Consider enabling them by passing, e.g., `--features="test_utils"` ``` --- [File an issue](https://github.com/adriangb/datafusion-benchmarking/issues) against this benchmark runner -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] - To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
Re: [PR] refactor: Update SortMergeJoin to use async spill abstractions [datafusion]
kumarUjjawal commented on PR #22230: URL: https://github.com/apache/datafusion/pull/22230#issuecomment-4619877677 run benchmark sort_merge_join -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] - To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
Re: [PR] refactor: Update SortMergeJoin to use async spill abstractions [datafusion]
2010YOUY01 commented on PR #22230: URL: https://github.com/apache/datafusion/pull/22230#issuecomment-4619750246 RE the ping β this PR is on my review queue, but I canβt promise exactly when Iβll be able to get to it :( Sorry for the delay, please proceed if anyone else can approve. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] - To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
Re: [PR] refactor: Update SortMergeJoin to use async spill abstractions [datafusion]
kumarUjjawal commented on PR #22230: URL: https://github.com/apache/datafusion/pull/22230#issuecomment-4542861684 @mbutrovich do you have bandwidth to look at this? -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] - To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
Re: [PR] refactor: Update SortMergeJoin to use async spill abstractions [datafusion]
pantShrey commented on code in PR #22230:
URL: https://github.com/apache/datafusion/pull/22230#discussion_r3296789577
##
datafusion/physical-plan/src/joins/sort_merge_join/bitwise_stream.rs:
##
@@ -785,24 +798,44 @@ impl BitwiseSortMergeJoinStream {
)
.count_ones();
-// Process spilled inner batches first (read back from disk).
-if let Some(spill_file) = &self.inner_key_spill {
-let file = BufReader::new(File::open(spill_file.path())?);
-let reader = StreamReader::try_new(file, None)?;
-for batch_result in reader {
-let inner_slice = batch_result?;
-matched_count = eval_filter_for_inner_slice(
-self.outer_is_left,
-filter,
-&outer_slice,
-&inner_slice,
-&mut self.matched,
-self.outer_offset,
-outer_group_len,
-matched_count,
-)?;
-if matched_count == outer_group_len {
-break;
+// Process spilled inner batches first asynchronously.
+if self.inner_key_spill.is_some() || self.spill_stream.is_some() {
Review Comment:
Sorry for missing this! I have added the `matched_count < outer_group_len`
guard to wrap the stream creation as well.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
-
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]
Re: [PR] refactor: Update SortMergeJoin to use async spill abstractions [datafusion]
pantShrey commented on code in PR #22230:
URL: https://github.com/apache/datafusion/pull/22230#discussion_r3296789577
##
datafusion/physical-plan/src/joins/sort_merge_join/bitwise_stream.rs:
##
@@ -785,24 +798,44 @@ impl BitwiseSortMergeJoinStream {
)
.count_ones();
-// Process spilled inner batches first (read back from disk).
-if let Some(spill_file) = &self.inner_key_spill {
-let file = BufReader::new(File::open(spill_file.path())?);
-let reader = StreamReader::try_new(file, None)?;
-for batch_result in reader {
-let inner_slice = batch_result?;
-matched_count = eval_filter_for_inner_slice(
-self.outer_is_left,
-filter,
-&outer_slice,
-&inner_slice,
-&mut self.matched,
-self.outer_offset,
-outer_group_len,
-matched_count,
-)?;
-if matched_count == outer_group_len {
-break;
+// Process spilled inner batches first asynchronously.
+if self.inner_key_spill.is_some() || self.spill_stream.is_some() {
Review Comment:
Sorry for missing this! I have moved the `matched_count < outer_group_len`
guard to wrap the stream creation as well.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
-
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]
Re: [PR] refactor: Update SortMergeJoin to use async spill abstractions [datafusion]
pantShrey commented on code in PR #22230:
URL: https://github.com/apache/datafusion/pull/22230#discussion_r3296743326
##
datafusion/physical-plan/src/joins/sort_merge_join/materializing_stream.rs:
##
@@ -997,6 +1112,7 @@ impl MaterializingSortMergeJoinStream {
.unwrap(); // Operation only return None if no
batches are spilled, here we ensure that at least one batch is spilled
buffered_batch.batch =
BufferedBatchState::Spilled(spill_file);
+self.spilled_batch_count += 1;
Review Comment:
Good catch! Added a decrement in the dequeue path for batches that are
popped while still in `Spilled` state.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
-
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]
Re: [PR] refactor: Update SortMergeJoin to use async spill abstractions [datafusion]
kumarUjjawal commented on code in PR #22230:
URL: https://github.com/apache/datafusion/pull/22230#discussion_r3296333431
##
datafusion/physical-plan/src/joins/sort_merge_join/materializing_stream.rs:
##
@@ -997,6 +1112,7 @@ impl MaterializingSortMergeJoinStream {
.unwrap(); // Operation only return None if no
batches are spilled, here we ensure that at least one batch is spilled
buffered_batch.batch =
BufferedBatchState::Spilled(spill_file);
+self.spilled_batch_count += 1;
Review Comment:
this also needs a decrement in the dequeue path.
##
datafusion/physical-plan/src/joins/sort_merge_join/bitwise_stream.rs:
##
@@ -785,24 +798,44 @@ impl BitwiseSortMergeJoinStream {
)
.count_ones();
-// Process spilled inner batches first (read back from disk).
-if let Some(spill_file) = &self.inner_key_spill {
-let file = BufReader::new(File::open(spill_file.path())?);
-let reader = StreamReader::try_new(file, None)?;
-for batch_result in reader {
-let inner_slice = batch_result?;
-matched_count = eval_filter_for_inner_slice(
-self.outer_is_left,
-filter,
-&outer_slice,
-&inner_slice,
-&mut self.matched,
-self.outer_offset,
-outer_group_len,
-matched_count,
-)?;
-if matched_count == outer_group_len {
-break;
+// Process spilled inner batches first asynchronously.
+if self.inner_key_spill.is_some() || self.spill_stream.is_some() {
Review Comment:
we should guard the stream creation
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
-
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]
Re: [PR] refactor: Update SortMergeJoin to use async spill abstractions [datafusion]
pantShrey commented on code in PR #22230:
URL: https://github.com/apache/datafusion/pull/22230#discussion_r3287038587
##
datafusion/physical-plan/src/joins/sort_merge_join/materializing_stream.rs:
##
@@ -917,6 +943,81 @@ impl MaterializingSortMergeJoinStream {
Poll::Pending
}
+/// Identifies which buffered batches are needed for the upcoming freeze
operation
+fn get_required_batch_indices(&self, buffered_freeze_count: usize) ->
Vec {
Review Comment:
sure, added `spilled_batch_count: usize` to track spill state . The fast
path now returns immediately when nothing is spilled, no scan needed.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
-
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]
Re: [PR] refactor: Update SortMergeJoin to use async spill abstractions [datafusion]
kumarUjjawal commented on code in PR #22230:
URL: https://github.com/apache/datafusion/pull/22230#discussion_r3285943557
##
datafusion/physical-plan/src/joins/sort_merge_join/materializing_stream.rs:
##
@@ -917,6 +943,81 @@ impl MaterializingSortMergeJoinStream {
Poll::Pending
}
+/// Identifies which buffered batches are needed for the upcoming freeze
operation
+fn get_required_batch_indices(&self, buffered_freeze_count: usize) ->
Vec {
Review Comment:
this gets rebuilt on every Polling transition? Could we short-circuit when
buffered_data.batches contains no Spilled variants? Avoids the O(n) scan on the
hot in-memory path.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
-
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]
Re: [PR] refactor: Update SortMergeJoin to use async spill abstractions [datafusion]
pantShrey commented on code in PR #22230:
URL: https://github.com/apache/datafusion/pull/22230#discussion_r3281058894
##
datafusion/physical-plan/src/joins/sort_merge_join/materializing_stream.rs:
##
@@ -917,6 +943,81 @@ impl MaterializingSortMergeJoinStream {
Poll::Pending
}
+/// Identifies which buffered batches are needed for the upcoming freeze
operation
+fn get_required_batch_indices(&self, buffered_freeze_count: usize) ->
Vec {
+let mut needed = vec![];
+
+// We need all batches that matched with streamed rows
+for chunk in &self.streamed_batch.output_indices {
+if let Some(idx) = chunk.buffered_batch_idx {
+needed.push(idx);
+}
+}
+
+// Full Joins need to emit null-joined rows, so we need batches up to
freeze_count
+if self.join_type == JoinType::Full {
+needed.extend(0..buffered_freeze_count);
+}
+
+needed.sort_unstable();
+needed.dedup();
+needed
+}
+
+/// Asynchronously reads spilled batches back into memory.
+/// Only processes the required indices to avoid OOMs.
+fn poll_spilled_batches(
+&mut self,
+cx: &mut Context<'_>,
+required_indices: &[usize],
+) -> Poll> {
+for &idx in required_indices {
+// Guard against indices that might be out of bounds if the queue
was cleared
+if idx >= self.buffered_data.batches.len() {
+continue;
+}
+
+let bb = &mut self.buffered_data.batches[idx];
+
+if let BufferedBatchState::Spilled(spill_file) = &bb.batch {
+if self.spill_stream.is_none() {
+let stream = self
+.spill_manager
+.read_spill_as_stream(spill_file.clone(), None)?;
+self.spill_stream = Some(stream);
+}
+
+match
ready!(self.spill_stream.as_mut().unwrap().poll_next_unpin(cx)) {
Review Comment:
I completely agree this needs coverage, I tried to build a mock for this
locally, but the challenge is that currently `read_spill_as_stream`
instantiates a real file stream internally, leaving no clean injection point.
This would become much cleaner once #21882 lands, at that point a
`MockSpillFile` can inject `Poll::Pending `cleanly without touching internals.
Would it be alright to handle this test as a follow-up once that PR lands?
Happy to open a tracking issue now if that's preferred.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
-
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]
Re: [PR] refactor: Update SortMergeJoin to use async spill abstractions [datafusion]
pantShrey commented on code in PR #22230:
URL: https://github.com/apache/datafusion/pull/22230#discussion_r3281058894
##
datafusion/physical-plan/src/joins/sort_merge_join/materializing_stream.rs:
##
@@ -917,6 +943,81 @@ impl MaterializingSortMergeJoinStream {
Poll::Pending
}
+/// Identifies which buffered batches are needed for the upcoming freeze
operation
+fn get_required_batch_indices(&self, buffered_freeze_count: usize) ->
Vec {
+let mut needed = vec![];
+
+// We need all batches that matched with streamed rows
+for chunk in &self.streamed_batch.output_indices {
+if let Some(idx) = chunk.buffered_batch_idx {
+needed.push(idx);
+}
+}
+
+// Full Joins need to emit null-joined rows, so we need batches up to
freeze_count
+if self.join_type == JoinType::Full {
+needed.extend(0..buffered_freeze_count);
+}
+
+needed.sort_unstable();
+needed.dedup();
+needed
+}
+
+/// Asynchronously reads spilled batches back into memory.
+/// Only processes the required indices to avoid OOMs.
+fn poll_spilled_batches(
+&mut self,
+cx: &mut Context<'_>,
+required_indices: &[usize],
+) -> Poll> {
+for &idx in required_indices {
+// Guard against indices that might be out of bounds if the queue
was cleared
+if idx >= self.buffered_data.batches.len() {
+continue;
+}
+
+let bb = &mut self.buffered_data.batches[idx];
+
+if let BufferedBatchState::Spilled(spill_file) = &bb.batch {
+if self.spill_stream.is_none() {
+let stream = self
+.spill_manager
+.read_spill_as_stream(spill_file.clone(), None)?;
+self.spill_stream = Some(stream);
+}
+
+match
ready!(self.spill_stream.as_mut().unwrap().poll_next_unpin(cx)) {
Review Comment:
I completely agree this needs coverage, I tried to build a mock for this
locally, but the challenge is that `read_spill_as_stream` instantiates a real
file stream internally, leaving no clean injection point. This would become
much cleaner once #21882 lands, at that point a `MockSpillFile` can inject
`Poll::Pending `cleanly without touching internals. Would it be alright to
handle this test as a follow-up once that PR lands? Happy to open a tracking
issue now if that's preferred.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
-
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]
Re: [PR] refactor: Update SortMergeJoin to use async spill abstractions [datafusion]
pantShrey commented on PR #22230: URL: https://github.com/apache/datafusion/pull/22230#issuecomment-4508186565 @kumarUjjawal Thanks so much for the review! I've pushed a new commit to address your feedbacks -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] - To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
Re: [PR] refactor: Update SortMergeJoin to use async spill abstractions [datafusion]
pantShrey commented on code in PR #22230:
URL: https://github.com/apache/datafusion/pull/22230#discussion_r3281058894
##
datafusion/physical-plan/src/joins/sort_merge_join/materializing_stream.rs:
##
@@ -917,6 +943,81 @@ impl MaterializingSortMergeJoinStream {
Poll::Pending
}
+/// Identifies which buffered batches are needed for the upcoming freeze
operation
+fn get_required_batch_indices(&self, buffered_freeze_count: usize) ->
Vec {
+let mut needed = vec![];
+
+// We need all batches that matched with streamed rows
+for chunk in &self.streamed_batch.output_indices {
+if let Some(idx) = chunk.buffered_batch_idx {
+needed.push(idx);
+}
+}
+
+// Full Joins need to emit null-joined rows, so we need batches up to
freeze_count
+if self.join_type == JoinType::Full {
+needed.extend(0..buffered_freeze_count);
+}
+
+needed.sort_unstable();
+needed.dedup();
+needed
+}
+
+/// Asynchronously reads spilled batches back into memory.
+/// Only processes the required indices to avoid OOMs.
+fn poll_spilled_batches(
+&mut self,
+cx: &mut Context<'_>,
+required_indices: &[usize],
+) -> Poll> {
+for &idx in required_indices {
+// Guard against indices that might be out of bounds if the queue
was cleared
+if idx >= self.buffered_data.batches.len() {
+continue;
+}
+
+let bb = &mut self.buffered_data.batches[idx];
+
+if let BufferedBatchState::Spilled(spill_file) = &bb.batch {
+if self.spill_stream.is_none() {
+let stream = self
+.spill_manager
+.read_spill_as_stream(spill_file.clone(), None)?;
+self.spill_stream = Some(stream);
+}
+
+match
ready!(self.spill_stream.as_mut().unwrap().poll_next_unpin(cx)) {
Review Comment:
I agree this needs coverage, I tried to build a mock for this locally, but
the challenge is that `read_spill_as_stream` instantiates a real file stream
internally, leaving no clean injection point. This would become much cleaner
once #21882 lands, at that point a `MockSpillFile` can inject `Poll::Pending
`cleanly without touching internals. Would it be alright to handle this test as
a follow-up once that PR lands? Happy to open a tracking issue now if that's
preferred.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
-
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]
Re: [PR] refactor: Update SortMergeJoin to use async spill abstractions [datafusion]
pantShrey commented on code in PR #22230:
URL: https://github.com/apache/datafusion/pull/22230#discussion_r3281014293
##
datafusion/physical-plan/src/joins/sort_merge_join/bitwise_stream.rs:
##
@@ -785,24 +794,40 @@ impl BitwiseSortMergeJoinStream {
)
.count_ones();
-// Process spilled inner batches first (read back from disk).
-if let Some(spill_file) = &self.inner_key_spill {
-let file = BufReader::new(File::open(spill_file.path())?);
-let reader = StreamReader::try_new(file, None)?;
-for batch_result in reader {
-let inner_slice = batch_result?;
-matched_count = eval_filter_for_inner_slice(
-self.outer_is_left,
-filter,
-&outer_slice,
-&inner_slice,
-&mut self.matched,
-self.outer_offset,
-outer_group_len,
-matched_count,
-)?;
-if matched_count == outer_group_len {
-break;
+// Process spilled inner batches first asynchronously.
+if self.inner_key_spill.is_some() || self.spill_stream.is_some() {
+if self.spill_stream.is_none()
+&& let Some(spill_file) = &self.inner_key_spill
+{
+let stream = self
+.spill_manager
+.read_spill_as_stream(spill_file.clone(), None)?;
+self.spill_stream = Some(stream);
+}
+
+while matched_count < outer_group_len {
+let stream = self.spill_stream.as_mut().unwrap();
+match ready!(stream.poll_next_unpin(cx)) {
+Some(Ok(inner_slice)) => {
+matched_count = eval_filter_for_inner_slice(
+self.outer_is_left,
+filter,
+&outer_slice,
+&inner_slice,
+&mut self.matched,
+self.outer_offset,
+outer_group_len,
+matched_count,
+)?;
+}
+Some(Err(e)) => {
+self.spill_stream = None;
+return Poll::Ready(Err(e));
+}
+None => {
Review Comment:
You are right. Although `None` is needed here to signify the normal end of
the stream, it should definitely have a guard to check for an unexpectedly
empty streams. I've added a `spill_stream_has_data: bool` flag to the struct.
`None` remains the normal EOF after reading batches, but it will now fire an
`internal_err!` if the very first poll returns `None`
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
-
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]
Re: [PR] refactor: Update SortMergeJoin to use async spill abstractions [datafusion]
pantShrey commented on code in PR #22230:
URL: https://github.com/apache/datafusion/pull/22230#discussion_r3280993470
##
datafusion/physical-plan/src/joins/sort_merge_join/materializing_stream.rs:
##
@@ -1603,29 +1707,18 @@ impl MaterializingSortMergeJoinStream {
let num_right_cols = self.buffered_schema.fields().len();
// Read each source batch once (spilled batches require disk I/O).
-// Track memory for each spilled batch at the point of deserialization
-// so the pool reflects actual usage as it grows.
-let spill_reservation = self.reservation.new_empty();
-let mut source_data: Vec> =
-Vec::with_capacity(source_batches.len());
-for &idx in &source_batches {
-let bb = &self.buffered_data.batches[idx];
-match &bb.batch {
-BufferedBatchState::InMemory(batch) => {
-source_data.push(Some(batch.clone()));
-}
-BufferedBatchState::Spilled(spill_file) => {
-spill_reservation.grow(bb.size_estimation);
-self.join_metrics
-.peak_mem_used()
-.set_max(self.reservation.size() +
spill_reservation.size());
-
-let file = BufReader::new(File::open(spill_file.path())?);
-let reader = StreamReader::try_new(file, None)?;
-source_data.push(reader.into_iter().next().transpose()?);
+let source_data: Vec = source_batches
+.iter()
+.map(|&idx| {
+let bb = &self.buffered_data.batches[idx];
+match &bb.batch {
+BufferedBatchState::InMemory(batch) => batch.clone(),
+BufferedBatchState::Spilled(_) => {
+unreachable!("Batches were unspilled")
Review Comment:
Fixed, replaced with `internal_err!` to fail the query gracefully instead of
panicking the executor.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
-
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]
Re: [PR] refactor: Update SortMergeJoin to use async spill abstractions [datafusion]
adriangbot commented on PR #22230: URL: https://github.com/apache/datafusion/pull/22230#issuecomment-4500378818 π€ Benchmark completed (GKE) | [trigger](https://github.com/apache/datafusion/pull/22230#issuecomment-4500228268) **Instance:** `c4a-highmem-16` (12 vCPU / 65 GiB) 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 ``` Comparing HEAD and smj-async-spill Benchmark smj.json βββββ³ββββ³βββ³ββββ β Query β HEAD β smj-async-spill βChange β β‘βββββββββββββββ© β QQuery 1 β 8.84 / 9.04 Β±0.18 / 9.39 ms β 8.55 / 8.91 Β±0.22 / 9.22 ms β no change β β QQuery 2 β 172.73 / 180.69 Β±4.46 / 186.23 ms β175.25 / 184.31 Β±4.73 / 188.20 ms β no change β β QQuery 3 β 104.29 / 106.12 Β±1.07 / 107.38 ms β108.85 / 109.88 Β±1.09 / 111.91 ms β no change β β QQuery 4 β28.16 / 28.30 Β±0.08 / 28.40 ms β 28.02 / 28.10 Β±0.08 / 28.25 ms β no change β β QQuery 5 β21.38 / 21.58 Β±0.17 / 21.87 ms β 21.72 / 22.01 Β±0.22 / 22.38 ms β no change β β QQuery 6 β 172.84 / 176.12 Β±2.21 / 179.28 ms β171.16 / 173.74 Β±2.09 / 177.35 ms β no change β β QQuery 7 β 208.09 / 210.38 Β±1.46 / 212.58 ms β211.94 / 217.30 Β±8.84 / 234.90 ms β no change β β QQuery 8 β19.96 / 20.33 Β±0.25 / 20.64 ms β 20.60 / 20.77 Β±0.13 / 20.91 ms β no change β β QQuery 9 β 216.06 / 220.14 Β±6.05 / 232.17 ms β219.13 / 221.85 Β±2.23 / 225.32 ms β no change β β QQuery 10 β69.04 / 72.78 Β±2.40 / 76.31 ms β 69.03 / 72.21 Β±2.20 / 75.86 ms β no change β β QQuery 11 β26.74 / 26.98 Β±0.15 / 27.18 ms β 26.82 / 26.93 Β±0.12 / 27.10 ms β no change β β QQuery 12 β66.81 / 69.19 Β±1.59 / 70.98 ms β 65.41 / 68.82 Β±1.73 / 70.11 ms β no change β β QQuery 13 β 98.51 / 102.79 Β±4.12 / 108.42 ms β101.84 / 108.00 Β±3.84 / 112.72 ms β 1.05x slower β β QQuery 14 β67.62 / 68.90 Β±0.87 / 70.26 ms β 68.23 / 69.30 Β±1.05 / 70.88 ms β no change β β QQuery 15 β67.94 / 71.23 Β±3.17 / 75.99 ms β 68.79 / 71.56 Β±2.72 / 76.65 ms β no change β β QQuery 16 β12.62 / 13.86 Β±2.20 / 18.26 ms
Re: [PR] refactor: Update SortMergeJoin to use async spill abstractions [datafusion]
adriangbot commented on PR #22230: URL: https://github.com/apache/datafusion/pull/22230#issuecomment-4500258025 Benchmark for [this request](https://github.com/apache/datafusion/pull/22230#issuecomment-4500230031) failed. Last 20 lines of output: Click to expand ``` Cloning into '/workspace/datafusion-branch'... From https://github.com/apache/datafusion * [new ref] refs/pull/22230/head -> smj-async-spill * branchmain -> FETCH_HEAD Switched to branch 'smj-async-spill' c8b784a01f5d0bcbe0dac806730fb61afc0be8ef Cloning into '/workspace/datafusion-base'... HEAD is now at c8b784a refactor(parquet-datasource): split sink and schema_coercion out of file_format.rs (#22347) rustc 1.95.0 (59807616e 2026-04-14) e3e1454a62315f24069c6f168b1515793546f60f c8b784a01f5d0bcbe0dac806730fb61afc0be8ef Blocking waiting for file lock on package cache Blocking waiting for file lock on package cache Blocking waiting for file lock on package cache error: target `sort_merge_join` in package `datafusion-physical-plan` requires the features: `test_utils` Consider enabling them by passing, e.g., `--features="test_utils"` ``` --- [File an issue](https://github.com/adriangb/datafusion-benchmarking/issues) against this benchmark runner -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] - To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
Re: [PR] refactor: Update SortMergeJoin to use async spill abstractions [datafusion]
adriangbot commented on PR #22230: URL: https://github.com/apache/datafusion/pull/22230#issuecomment-4500257111 π€ Criterion benchmark running (GKE) | [trigger](https://github.com/apache/datafusion/pull/22230#issuecomment-4500230031) **Instance:** `c4a-highmem-16` (12 vCPU / 65 GiB) | `Linux bench-c4500230031-223-g944f 6.12.68+ #1 SMP Wed Apr 1 02:23:28 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 smj-async-spill (e3e1454a62315f24069c6f168b1515793546f60f) to c8b784a (merge-base) [diff](https://github.com/apache/datafusion/compare/c8b784a01f5d0bcbe0dac806730fb61afc0be8ef..e3e1454a62315f24069c6f168b1515793546f60f) BENCH_NAME=sort_merge_join BENCH_COMMAND=cargo bench --features=parquet --bench sort_merge_join BENCH_FILTER= Results will be posted here when complete --- [File an issue](https://github.com/adriangb/datafusion-benchmarking/issues) against this benchmark runner -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] - To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
Re: [PR] refactor: Update SortMergeJoin to use async spill abstractions [datafusion]
adriangbot commented on PR #22230: URL: https://github.com/apache/datafusion/pull/22230#issuecomment-4500253718 π€ Benchmark running (GKE) | [trigger](https://github.com/apache/datafusion/pull/22230#issuecomment-4500228268) **Instance:** `c4a-highmem-16` (12 vCPU / 65 GiB) | `Linux bench-c4500228268-222-7db4b 6.12.68+ #1 SMP Wed Apr 1 02:23:28 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 smj-async-spill (e3e1454a62315f24069c6f168b1515793546f60f) to c8b784a (merge-base) [diff](https://github.com/apache/datafusion/compare/c8b784a01f5d0bcbe0dac806730fb61afc0be8ef..e3e1454a62315f24069c6f168b1515793546f60f) using: smj Results will be posted here when complete --- [File an issue](https://github.com/adriangb/datafusion-benchmarking/issues) against this benchmark runner -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] - To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
Re: [PR] refactor: Update SortMergeJoin to use async spill abstractions [datafusion]
alamb commented on PR #22230: URL: https://github.com/apache/datafusion/pull/22230#issuecomment-4500230031 run benchmark sort_merge_join -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] - To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
Re: [PR] refactor: Update SortMergeJoin to use async spill abstractions [datafusion]
alamb commented on PR #22230: URL: https://github.com/apache/datafusion/pull/22230#issuecomment-4500228268 run benchmarks smj -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] - To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
Re: [PR] refactor: Update SortMergeJoin to use async spill abstractions [datafusion]
alamb commented on PR #22230: URL: https://github.com/apache/datafusion/pull/22230#issuecomment-4500112072 @mbutrovich or @2010YOUY01 wonder if you have time to review this PR (at least give us some hints about how to run the relevant performance tests for sort merge join) to make sure we aren't messing up erformance The larger context is to try and make it easier to extend spilling (so not always to local files) -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] - To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
Re: [PR] refactor: Update SortMergeJoin to use async spill abstractions [datafusion]
kumarUjjawal commented on code in PR #22230:
URL: https://github.com/apache/datafusion/pull/22230#discussion_r3275199037
##
datafusion/physical-plan/src/joins/sort_merge_join/bitwise_stream.rs:
##
@@ -785,24 +794,40 @@ impl BitwiseSortMergeJoinStream {
)
.count_ones();
-// Process spilled inner batches first (read back from disk).
-if let Some(spill_file) = &self.inner_key_spill {
-let file = BufReader::new(File::open(spill_file.path())?);
-let reader = StreamReader::try_new(file, None)?;
-for batch_result in reader {
-let inner_slice = batch_result?;
-matched_count = eval_filter_for_inner_slice(
-self.outer_is_left,
-filter,
-&outer_slice,
-&inner_slice,
-&mut self.matched,
-self.outer_offset,
-outer_group_len,
-matched_count,
-)?;
-if matched_count == outer_group_len {
-break;
+// Process spilled inner batches first asynchronously.
+if self.inner_key_spill.is_some() || self.spill_stream.is_some() {
+if self.spill_stream.is_none()
+&& let Some(spill_file) = &self.inner_key_spill
+{
+let stream = self
+.spill_manager
+.read_spill_as_stream(spill_file.clone(), None)?;
+self.spill_stream = Some(stream);
+}
+
+while matched_count < outer_group_len {
+let stream = self.spill_stream.as_mut().unwrap();
+match ready!(stream.poll_next_unpin(cx)) {
+Some(Ok(inner_slice)) => {
+matched_count = eval_filter_for_inner_slice(
+self.outer_is_left,
+filter,
+&outer_slice,
+&inner_slice,
+&mut self.matched,
+self.outer_offset,
+outer_group_len,
+matched_count,
+)?;
+}
+Some(Err(e)) => {
+self.spill_stream = None;
+return Poll::Ready(Err(e));
+}
+None => {
Review Comment:
Can we make the empty-spill invariant consistent here? In the other async
spill-read path, None becomes internal_err!("Spill file was empty"), but here
we treat None as normal EOF and just break. Since these spill files come from
non-empty buffered data, I think this should either error here as well, or we
should document why empty spill streams are valid in this path.
##
datafusion/physical-plan/src/joins/sort_merge_join/materializing_stream.rs:
##
@@ -1603,29 +1707,18 @@ impl MaterializingSortMergeJoinStream {
let num_right_cols = self.buffered_schema.fields().len();
// Read each source batch once (spilled batches require disk I/O).
-// Track memory for each spilled batch at the point of deserialization
-// so the pool reflects actual usage as it grows.
-let spill_reservation = self.reservation.new_empty();
-let mut source_data: Vec> =
-Vec::with_capacity(source_batches.len());
-for &idx in &source_batches {
-let bb = &self.buffered_data.batches[idx];
-match &bb.batch {
-BufferedBatchState::InMemory(batch) => {
-source_data.push(Some(batch.clone()));
-}
-BufferedBatchState::Spilled(spill_file) => {
-spill_reservation.grow(bb.size_estimation);
-self.join_metrics
-.peak_mem_used()
-.set_max(self.reservation.size() +
spill_reservation.size());
-
-let file = BufReader::new(File::open(spill_file.path())?);
-let reader = StreamReader::try_new(file, None)?;
-source_data.push(reader.into_iter().next().transpose()?);
+let source_data: Vec = source_batches
+.iter()
+.map(|&idx| {
+let bb = &self.buffered_data.batches[idx];
+match &bb.batch {
+BufferedBatchState::InMemory(batch) => batch.clone(),
+BufferedBatchState::Spilled(_) => {
+unreachable!("Batches were unspilled")
Review Comment:
Iβd prefer `internal_err!` here instead of unreachable!. The same "should
have been unspilled already" precondition in
`fetch_right_columns_from_batch_by_idxs` already returns a regular internal
error. If a future caller misses the unspill step, failing the query is much
better than panicking
Re: [PR] refactor: Update SortMergeJoin to use async spill abstractions [datafusion]
pantShrey commented on PR #22230: URL: https://github.com/apache/datafusion/pull/22230#issuecomment-4496494148 Fixed the failing CI test fuzz_cases::join_fuzz::test_filtered_join_spill_fuzz -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] - To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
