gene-bordegaray commented on code in PR #24417:
URL: https://github.com/apache/datafusion/pull/24417#discussion_r3797764270


##########
datafusion/physical-plan/src/joins/hash_join/shared_bounds.rs:
##########
@@ -611,188 +611,178 @@ impl SharedBuildAccumulator {
 
     fn build_filter(&self, finalize_input: FinalizeInput) -> Result<()> {
         match finalize_input {
-            FinalizeInput::CollectLeft(partition) => match partition {
-                PartitionStatus::Reported(partition_data) => {
+            FinalizeInput::CollectLeft(partition) => {
+                self.build_collect_left_filter(partition)
+            }
+            FinalizeInput::Partitioned(partitions) => {
+                self.build_partitioned_filter(&partitions)

Review Comment:
   sonce we own here we could consume this vec to eliminate some clones I 
believe



##########
datafusion/physical-plan/src/joins/hash_join/shared_bounds.rs:
##########
@@ -611,188 +611,178 @@ impl SharedBuildAccumulator {
 
     fn build_filter(&self, finalize_input: FinalizeInput) -> Result<()> {
         match finalize_input {
-            FinalizeInput::CollectLeft(partition) => match partition {
-                PartitionStatus::Reported(partition_data) => {
+            FinalizeInput::CollectLeft(partition) => {
+                self.build_collect_left_filter(partition)
+            }
+            FinalizeInput::Partitioned(partitions) => {
+                self.build_partitioned_filter(&partitions)
+            }
+        }
+    }
+
+    /// Builds the single global filter used by a collect-left join.
+    fn build_collect_left_filter(&self, partition: PartitionStatus) -> 
Result<()> {
+        match partition {
+            PartitionStatus::Reported(partition_data) => {
+                let membership_expr = create_membership_predicate(
+                    &self.on_right,
+                    partition_data.pushdown.clone(),
+                    &HASH_JOIN_SEED,
+                    self.probe_schema.as_ref(),
+                )?;
+                let bounds_expr =
+                    create_bounds_predicate(&self.on_right, 
&partition_data.bounds);
+
+                if let Some(filter_expr) =
+                    combine_membership_and_bounds(membership_expr, bounds_expr)
+                {
+                    self.dynamic_filter.update(self.preserve_probe_nulls(
+                        filter_expr,
+                        partition_data.keys_have_null,
+                    )?)?;
+                }
+                Ok(())
+            }
+            PartitionStatus::Pending => datafusion_common::internal_err!(
+                "attempted to finalize collect-left dynamic filter without 
reported build data"
+            ),
+            PartitionStatus::CanceledUnknown => 
datafusion_common::internal_err!(
+                "collect-left dynamic filter cannot finalize with canceled 
build data"
+            ),
+        }
+    }
+
+    fn build_partitioned_filter(&self, partitions: &[PartitionStatus]) -> 
Result<()> {

Review Comment:
   can we also add a doc comment



##########
datafusion/physical-plan/src/joins/hash_join/shared_bounds.rs:
##########
@@ -611,188 +611,178 @@ impl SharedBuildAccumulator {
 
     fn build_filter(&self, finalize_input: FinalizeInput) -> Result<()> {
         match finalize_input {
-            FinalizeInput::CollectLeft(partition) => match partition {
-                PartitionStatus::Reported(partition_data) => {
+            FinalizeInput::CollectLeft(partition) => {
+                self.build_collect_left_filter(partition)
+            }
+            FinalizeInput::Partitioned(partitions) => {
+                self.build_partitioned_filter(&partitions)
+            }
+        }
+    }
+
+    /// Builds the single global filter used by a collect-left join.
+    fn build_collect_left_filter(&self, partition: PartitionStatus) -> 
Result<()> {
+        match partition {
+            PartitionStatus::Reported(partition_data) => {
+                let membership_expr = create_membership_predicate(
+                    &self.on_right,
+                    partition_data.pushdown.clone(),
+                    &HASH_JOIN_SEED,
+                    self.probe_schema.as_ref(),
+                )?;
+                let bounds_expr =
+                    create_bounds_predicate(&self.on_right, 
&partition_data.bounds);
+
+                if let Some(filter_expr) =
+                    combine_membership_and_bounds(membership_expr, bounds_expr)
+                {
+                    self.dynamic_filter.update(self.preserve_probe_nulls(
+                        filter_expr,
+                        partition_data.keys_have_null,
+                    )?)?;
+                }
+                Ok(())
+            }
+            PartitionStatus::Pending => datafusion_common::internal_err!(
+                "attempted to finalize collect-left dynamic filter without 
reported build data"
+            ),
+            PartitionStatus::CanceledUnknown => 
datafusion_common::internal_err!(
+                "collect-left dynamic filter cannot finalize with canceled 
build data"
+            ),
+        }
+    }
+
+    fn build_partitioned_filter(&self, partitions: &[PartitionStatus]) -> 
Result<()> {
+        let mut partition_filters = Vec::with_capacity(partitions.len());
+        let mut real_partition_ids = Vec::new();
+        let mut empty_partition_ids = Vec::new();
+        let mut has_canceled_unknown = false;
+        let mut keys_have_null = false;
+
+        for (partition_id, partition) in partitions.iter().enumerate() {
+            match partition {
+                PartitionStatus::Reported(partition)
+                    if matches!(partition.pushdown, PushdownStrategy::Empty) =>
+                {
+                    empty_partition_ids.push(partition_id);
+                    partition_filters.push(lit(false));
+                }
+                PartitionStatus::Reported(partition) => {
+                    real_partition_ids.push(partition_id);
+                    keys_have_null |= partition.keys_have_null;
                     let membership_expr = create_membership_predicate(
                         &self.on_right,
-                        partition_data.pushdown.clone(),
+                        partition.pushdown.clone(),
                         &HASH_JOIN_SEED,
                         self.probe_schema.as_ref(),
                     )?;
                     let bounds_expr =
-                        create_bounds_predicate(&self.on_right, 
&partition_data.bounds);
-
-                    if let Some(filter_expr) =
+                        create_bounds_predicate(&self.on_right, 
&partition.bounds);
+                    let then_expr =
                         combine_membership_and_bounds(membership_expr, 
bounds_expr)
-                    {
-                        self.dynamic_filter.update(self.preserve_probe_nulls(
-                            filter_expr,
-                            partition_data.keys_have_null,
-                        )?)?;
-                    }
-                }
-                PartitionStatus::Pending => {
-                    return datafusion_common::internal_err!(
-                        "attempted to finalize collect-left dynamic filter 
without reported build data"
-                    );
+                            .unwrap_or_else(|| lit(true));
+                    partition_filters.push(then_expr);
                 }
                 PartitionStatus::CanceledUnknown => {
+                    has_canceled_unknown = true;
+                    partition_filters.push(lit(true));
+                    // A canceled partition's build content is unknown, so it
+                    // may hold a NULL key.
+                    keys_have_null = true;
+                }
+                PartitionStatus::Pending => {
                     return datafusion_common::internal_err!(
-                        "collect-left dynamic filter cannot finalize with 
canceled build data"
+                        "attempted to finalize dynamic filter with pending 
partition"
                     );
                 }
-            },
-            FinalizeInput::Partitioned(partitions) => {
-                let num_partitions = partitions.len();
-                let mut partition_filters = Vec::with_capacity(num_partitions);
-                let mut real_partition_ids = Vec::new();
-                let mut empty_partition_ids = Vec::new();
-                let mut has_canceled_unknown = false;
-                let mut keys_have_null = false;
-
-                for (partition_id, partition) in partitions.iter().enumerate() 
{
-                    match partition {
-                        PartitionStatus::Reported(partition)
-                            if matches!(partition.pushdown, 
PushdownStrategy::Empty) =>
-                        {
-                            empty_partition_ids.push(partition_id);
-                            partition_filters.push(lit(false));
-                        }
-                        PartitionStatus::Reported(partition) => {
-                            real_partition_ids.push(partition_id);
-                            keys_have_null |= partition.keys_have_null;
-                            let membership_expr = create_membership_predicate(
-                                &self.on_right,
-                                partition.pushdown.clone(),
-                                &HASH_JOIN_SEED,
-                                self.probe_schema.as_ref(),
-                            )?;
-                            let bounds_expr = create_bounds_predicate(
-                                &self.on_right,
-                                &partition.bounds,
-                            );
-                            let then_expr = combine_membership_and_bounds(
-                                membership_expr,
-                                bounds_expr,
-                            )
-                            .unwrap_or_else(|| lit(true));
-                            partition_filters.push(then_expr);
-                        }
-                        PartitionStatus::CanceledUnknown => {
-                            has_canceled_unknown = true;
-                            partition_filters.push(lit(true));
-                            // A canceled partition's build content is 
unknown, so it
-                            // may hold a NULL key.
-                            keys_have_null = true;
-                        }
-                        PartitionStatus::Pending => {
-                            return datafusion_common::internal_err!(
-                                "attempted to finalize dynamic filter with 
pending partition"
-                            );
-                        }
-                    }
-                }
-
-                let filter_expr = if has_canceled_unknown
-                    && real_partition_ids.is_empty()
-                    && empty_partition_ids.is_empty()
-                {
-                    lit(true)
-                } else if !has_canceled_unknown && 
real_partition_ids.is_empty() {
-                    lit(false)
-                } else if !has_canceled_unknown
-                    && real_partition_ids.len() == 1
-                    && empty_partition_ids.len() + 1 == num_partitions
-                {
-                    Arc::clone(&partition_filters[real_partition_ids[0]])
-                } else if let Some(range_partitioning) = 
&self.probe_range_partitioning {
-                    // Range partitioning
-                    assert_or_internal_err!(
-                        partition_filters.len() == 
range_partitioning.partition_count(),
-                        "Dynamic filter partition count {} does not match 
Range partition count {}",
-                        partition_filters.len(),
-                        range_partitioning.partition_count()
-                    );
-                    let routing_range_expr = Arc::new(RangeExpr::try_new(
-                        self.on_right.clone(),
-                        range_partitioning,
-                    )?)
-                        as Arc<dyn PhysicalExpr>;
-                    let else_expr = partition_filters
-                        .pop()
-                        .expect("Range partitioning always has at least one 
partition");
-
-                    // CASE range_partition(key)
-                    //   WHEN 0 THEN F0
-                    //   WHEN 1 THEN F1
-                    //   ...
-                    //   ELSE Fn
-                    // END
-                    let when_then_expr = partition_filters
-                        .into_iter()
-                        .enumerate()
-                        .map(|(partition_id, then_expr)| {
-                            (
-                                lit(ScalarValue::UInt64(Some(partition_id as 
u64))),
-                                then_expr,
-                            )
-                        })
-                        .collect();
-
-                    Arc::new(CaseExpr::try_new(
-                        Some(routing_range_expr),
-                        when_then_expr,
-                        Some(else_expr),
-                    )?) as Arc<dyn PhysicalExpr>
-                } else {
-                    // Hash partitioning
-                    let routing_hash_expr = Arc::new(HashExpr::new(
-                        self.on_right.clone(),
-                        self.repartition_random_state.clone(),
-                        "hash_repartition".to_string(),
-                    ))
-                        as Arc<dyn PhysicalExpr>;
-                    let modulo_expr = Arc::new(BinaryExpr::new(
-                        routing_hash_expr,
-                        Operator::Modulo,
-                        lit(ScalarValue::UInt64(Some(num_partitions as u64))),
-                    )) as Arc<dyn PhysicalExpr>;
-
-                    let mut when_then_branches = if has_canceled_unknown {
-                        empty_partition_ids
-                            .into_iter()
-                            .map(|partition_id| {
-                                (
-                                    lit(ScalarValue::UInt64(Some(partition_id 
as u64))),
-                                    lit(false),
-                                )
-                            })
-                            .collect::<Vec<_>>()
-                    } else {
-                        vec![]
-                    };
-                    
when_then_branches.extend(real_partition_ids.into_iter().map(
-                        |partition_id| {
-                            (
-                                lit(ScalarValue::UInt64(Some(partition_id as 
u64))),
-                                Arc::clone(&partition_filters[partition_id]),
-                            )
-                        },
-                    ));
-
-                    Arc::new(CaseExpr::try_new(
-                        Some(modulo_expr),
-                        when_then_branches,
-                        Some(lit(has_canceled_unknown)),
-                    )?) as Arc<dyn PhysicalExpr>
-                };
-
-                self.dynamic_filter
-                    .update(self.preserve_probe_nulls(filter_expr, 
keys_have_null)?)?;
             }
         }
 
-        Ok(())
+        let all_partitions_canceled = has_canceled_unknown
+            && real_partition_ids.is_empty()
+            && empty_partition_ids.is_empty();
+        let all_partitions_empty = !has_canceled_unknown && 
real_partition_ids.is_empty();
+        let one_non_empty_partition =
+            !has_canceled_unknown && real_partition_ids.len() == 1;
+
+        let filter_expr = if all_partitions_canceled {
+            // No build data is known, so filtering any probe row could 
discard a match.
+            lit(true)
+        } else if all_partitions_empty {
+            // No build row exists, so no probe row can match.
+            lit(false)
+        } else if one_non_empty_partition {
+            // Only one build partition contains rows, so its filter covers 
every
+            // possible probe match without routing.
+            Arc::clone(&partition_filters[real_partition_ids[0]])
+        } else {
+            // Builds the shared sparse `CASE` for partition filter routing.
+            // If no canceled partitions, `ELSE false` covers omitted empty 
partitions, and vice versa.

Review Comment:
   this is not that clear could we explicitly state expected behavior on 
cancelled partitions



##########
datafusion/physical-plan/src/joins/hash_join/shared_bounds.rs:
##########
@@ -611,188 +611,178 @@ impl SharedBuildAccumulator {
 
     fn build_filter(&self, finalize_input: FinalizeInput) -> Result<()> {
         match finalize_input {
-            FinalizeInput::CollectLeft(partition) => match partition {
-                PartitionStatus::Reported(partition_data) => {
+            FinalizeInput::CollectLeft(partition) => {
+                self.build_collect_left_filter(partition)
+            }
+            FinalizeInput::Partitioned(partitions) => {
+                self.build_partitioned_filter(&partitions)
+            }
+        }
+    }
+
+    /// Builds the single global filter used by a collect-left join.
+    fn build_collect_left_filter(&self, partition: PartitionStatus) -> 
Result<()> {
+        match partition {
+            PartitionStatus::Reported(partition_data) => {
+                let membership_expr = create_membership_predicate(
+                    &self.on_right,
+                    partition_data.pushdown.clone(),

Review Comment:
   same thing here with cloning



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