Re: [PR] refactor: Update SortMergeJoin to use async spill abstractions [datafusion]

2026-06-16 Thread via GitHub


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]

2026-06-16 Thread via GitHub


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]

2026-06-16 Thread via GitHub


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]

2026-06-15 Thread via GitHub


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]

2026-06-15 Thread via GitHub


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]

2026-06-15 Thread via GitHub


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]

2026-06-15 Thread via GitHub


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]

2026-06-10 Thread via GitHub


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]

2026-06-09 Thread via GitHub


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]

2026-06-09 Thread via GitHub


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]

2026-06-09 Thread via GitHub


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]

2026-06-09 Thread via GitHub


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]

2026-06-05 Thread via GitHub


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]

2026-06-04 Thread via GitHub


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]

2026-06-04 Thread via GitHub


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]

2026-06-04 Thread via GitHub


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]

2026-06-03 Thread via GitHub


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]

2026-05-26 Thread via GitHub


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]

2026-05-25 Thread via GitHub


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]

2026-05-25 Thread via GitHub


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]

2026-05-25 Thread via GitHub


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]

2026-05-24 Thread via GitHub


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]

2026-05-22 Thread via GitHub


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]

2026-05-21 Thread via GitHub


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]

2026-05-21 Thread via GitHub


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]

2026-05-21 Thread via GitHub


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]

2026-05-21 Thread via GitHub


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]

2026-05-21 Thread via GitHub


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]

2026-05-21 Thread via GitHub


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]

2026-05-21 Thread via GitHub


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]

2026-05-20 Thread via GitHub


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]

2026-05-20 Thread via GitHub


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]

2026-05-20 Thread via GitHub


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]

2026-05-20 Thread via GitHub


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]

2026-05-20 Thread via GitHub


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]

2026-05-20 Thread via GitHub


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]

2026-05-20 Thread via GitHub


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]

2026-05-20 Thread via GitHub


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]

2026-05-20 Thread via GitHub


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]