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


##########
datafusion/physical-plan/src/joins/hash_join/stream.rs:
##########
@@ -386,6 +388,9 @@ pub(super) struct HashJoinStream {
     probe_indices_buffer: Vec<u32>,
     /// Scratch space for build indices during hash lookup
     build_indices_buffer: Vec<u64>,
+    /// Key comparator for the current probe batch, built on first use and
+    /// reused by every chunk of that batch
+    probe_key_comparator: Option<JoinKeyComparator>,

Review Comment:
   The comparator is only valid for one probe batch, but it lives on the stream 
and relies on hand-placed `= None` resets (`fetch_probe_batch`, limit path). 
Any future path that installs a new `ProcessProbeBatchState` without a reset 
silently drops matches; removing the `fetch_probe_batch` reset fails 
`join_inner_multi_key_across_probe_batches`. Storing it on 
`ProcessProbeBatchState` ties its lifetime to the batch, removes both resets, 
and frees the probe arrays when the batch finishes instead of at the next 
fetch. I checked this locally: clippy is clean and the `joins::hash_join` tests 
pass.
   
   ```diff
   -#[derive(Debug, Clone)]
   +#[derive(Debug)]
    pub(super) enum HashJoinStreamState {
   ```
   
   ```diff
   -#[derive(Debug, Clone)]
   +#[derive(Debug)]
    pub(super) struct ProcessProbeBatchState {
        ...
        matched_probe_idx: Option<u32>,
   +    /// Key comparator for this batch, built on first use and reused by 
every
   +    /// chunk
   +    key_comparator: Option<JoinKeyComparator>,
    }
   ```
   
   ```diff
                            matched_probe_idx: None,
   +                        key_comparator: None,
                        });
   ```
   
   ```diff
   -                &mut self.probe_key_comparator,
   +                &mut state.key_comparator,
   ```
   
   Then drop the `probe_key_comparator` field, its initializer, and both 
`self.probe_key_comparator = None;` lines. In `utils.rs`:
   
   ```rs
   impl Debug for JoinKeyComparator {
       fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
           f.debug_struct("JoinKeyComparator").finish_non_exhaustive()
       }
   }
   ```



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