@@ -577,6 +577,12 @@ class VectorizedHashJoinOperator : public VectorizedOperator {
577577 std::vector<bool > left_matched_in_batch_;
578578 std::vector<size_t > unmatched_indices_;
579579
580+ // Probe state for resumable bucket scanning (prevents batch overflow)
581+ bool resuming_bucket_scan_ = false ; // True if we're resuming a mid-bucket scan
582+ size_t resumed_bucket_idx_ = 0 ; // Bucket index when resuming
583+ size_t resumed_entry_idx_ = 0 ; // Entry index within bucket when resuming
584+ common::Value resumed_key_val_; // Key value being probed when resuming
585+
580586 // Join type
581587 JoinType join_type_;
582588
@@ -711,10 +717,52 @@ class VectorizedHashJoinOperator : public VectorizedOperator {
711717 left_row_idx_ = 0 ;
712718 // Reset matched tracking for new batch
713719 std::fill (left_matched_in_batch_.begin (), left_matched_in_batch_.end (), false );
720+ // Clear resume state when advancing to new batch
721+ resuming_bucket_scan_ = false ;
714722 }
715723
716724 // Process rows in current batch
717725 while (left_row_idx_ < left_batch_->row_count () && out_batch.row_count () < BATCH_SIZE ) {
726+ // Check if we need to resume an interrupted bucket scan
727+ if (resuming_bucket_scan_) {
728+ // We were in the middle of scanning a bucket - resume from saved position
729+ const auto & key_val = resumed_key_val_;
730+ auto & bucket = buckets_[resumed_bucket_idx_];
731+ bool found_match = left_matched_in_batch_[left_row_idx_];
732+
733+ // Resume scanning bucket from resumed_entry_idx_
734+ for (size_t i = resumed_entry_idx_; i < bucket.key_values .size (); ++i) {
735+ if (out_batch.row_count () >= BATCH_SIZE ) {
736+ // Batch full - save state and return
737+ resuming_bucket_scan_ = true ;
738+ resumed_bucket_idx_ = resumed_bucket_idx_;
739+ resumed_entry_idx_ = i;
740+ resumed_key_val_ = key_val;
741+ return true ; // Caller must consume batch before continuing
742+ }
743+
744+ const auto & bucket_key = bucket.key_values [i][right_key_col_idx_];
745+ if (bucket_key == key_val) {
746+ emit_joined_row (out_batch, left_row_idx_, bucket.payload_rows [i]);
747+ found_match = true ;
748+ if (join_type_ == JoinType::Left) {
749+ left_matched_in_batch_[left_row_idx_] = true ;
750+ }
751+ }
752+ }
753+
754+ // Finished scanning this bucket
755+ resuming_bucket_scan_ = false ;
756+
757+ // Track unmatched for LEFT join
758+ if (join_type_ == JoinType::Left && !found_match) {
759+ unmatched_indices_.push_back (left_row_idx_);
760+ }
761+
762+ left_row_idx_++;
763+ continue ;
764+ }
765+
718766 const auto & key_val = left_batch_->get_column (left_key_col_idx_).get (left_row_idx_);
719767
720768 if (key_val.is_null ()) {
@@ -732,6 +780,15 @@ class VectorizedHashJoinOperator : public VectorizedOperator {
732780 // Search for match in this bucket
733781 bool found_match = false ;
734782 for (size_t i = 0 ; i < bucket.key_values .size (); ++i) {
783+ if (out_batch.row_count () >= BATCH_SIZE ) {
784+ // Batch full - save state and return
785+ resuming_bucket_scan_ = true ;
786+ resumed_bucket_idx_ = bucket_idx;
787+ resumed_entry_idx_ = i;
788+ resumed_key_val_ = key_val;
789+ return true ; // Caller must consume batch before continuing
790+ }
791+
735792 const auto & bucket_key = bucket.key_values [i][right_key_col_idx_];
736793 if (bucket_key == key_val) {
737794 // Match found - emit row
@@ -740,7 +797,7 @@ class VectorizedHashJoinOperator : public VectorizedOperator {
740797 if (join_type_ == JoinType::Left) {
741798 left_matched_in_batch_[left_row_idx_] = true ;
742799 }
743- break ; // Each left row matches at most one right row
800+ // Continue scanning bucket for all matching right rows
744801 }
745802 }
746803
0 commit comments