jayzhan211 commented on code in PR #25933:
URL: https://github.com/apache/datafusion/pull/25933#discussion_r4175545320


##########
datafusion/physical-plan/src/joins/sort_merge_join/bitwise_stream.rs:
##########
@@ -669,7 +669,11 @@ impl BitwiseSortMergeJoinStream {
 
             let inner_batch = self.inner_batch.as_ref().unwrap();
             let slice = inner_batch.slice(from, group_end - from);
-            self.inner_buffer_size += slice.get_array_memory_size();
+            // A slice reports its parent batch's full buffers, so charge only
+            // the rows the group holds. View arrays still count their parent's
+            // data buffers, and a group spanning an inner batch boundary keeps
+            // the earlier parent batch alive, so this can be one batch low.
+            self.inner_buffer_size += slice.get_sliced_size()?;

Review Comment:
   When `group_end == num_inner`, `next_inner_batch()` moves the cursor on, so 
the buffered slice becomes the only thing keeping its whole parent alive. It is 
charged only its rows (8192-row batch, group at row 8191: ~98 KB held, 12 B 
charged), whether or not the key continues. Main charged this case exactly. 
Charging the full parent there costs at most one spill per inner batch under 
pools smaller than a batch; I measured 2/2/1/1 spills in the new test, and the 
other 247 SMJ tests pass. Fine to handle in a follow-up.
   
   ```diff
   -            // A slice reports its parent batch's full buffers, so charge 
only
   -            // the rows the group holds. View arrays still count their 
parent's
   -            // data buffers, and a group spanning an inner batch boundary 
keeps
   -            // the earlier parent batch alive, so this can be one batch low.
   -            self.inner_buffer_size += slice.get_sliced_size()?;
   +            // A group ending inside the batch shares the current inner 
batch,
   +            // so charge only its rows. A group reaching the batch end keeps
   +            // the whole parent alive once the cursor advances, so charge 
all
   +            // of it. View arrays still count their parent's data buffers.
   +            self.inner_buffer_size += if group_end < num_inner {
   +                slice.get_sliced_size()?
   +            } else {
   +                slice.get_array_memory_size()
   +            };
   ```
   
   and in `bitwise_small_key_groups_charged_by_sliced_size`:
   
   ```rs
   let spill_count = metrics.spill_count().unwrap();
   assert!(
       spill_count <= NUM_BATCHES as usize,
       "{spill_count} spills under a {memory_limit} byte pool for {join_type:?}"
   );
   ```



-- 
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]

Reply via email to