viirya commented on code in PR #19857:
URL: https://github.com/apache/datafusion/pull/19857#discussion_r4109556476
##########
datafusion/physical-plan/src/joins/hash_join/stream.rs:
##########
@@ -1541,6 +1556,58 @@ fn mark_null_candidates_for_probe_batch(
Ok(())
}
+/// Keeps the candidate pairs whose multi-column `NOT IN` value tuples are not
a
+/// definite mismatch, i.e. every element pair is equal or involves a NULL.
+///
+/// Such a pair compares UNKNOWN when some element is NULL (TRUE pairs, with no
+/// NULL, are found by the hash lookup instead), whereas a pair with a definite
+/// mismatch compares FALSE whatever its NULLs are: `(NULL, 1) = (2, 3)` is
+/// `UNKNOWN AND FALSE`, which is FALSE.
+fn retain_value_mismatch_free(
+ build_value_keys: &[ArrayRef],
+ probe_value_keys: &[ArrayRef],
+ build_indices: UInt64Array,
+ probe_indices: UInt32Array,
+) -> Result<(UInt64Array, UInt32Array)> {
+ if build_indices.is_empty() {
+ return Ok((build_indices, probe_indices));
+ }
+ let comparators = build_value_keys
+ .iter()
+ .zip(probe_value_keys)
+ .map(|(build, probe)| {
+ Ok((
+ build.logical_nulls(),
+ probe.logical_nulls(),
+ // Same total order the hash join's own key comparison uses.
Review Comment:
Good catch, all three queries reproduced. Fixed in 8728b2c by taking your
batch-level suggestion: each chunk of candidate pairs is now compared with
`take` on both sides plus `datum::apply_cmp(Operator::Eq)`, so the UNKNOWN
check uses the same float-normalizing `=` as SQL (nested types included), and
the work is proportional to the chunk rather than a scan of the full build
arrays. Your three queries are now in `null_aware_multi_column.slt`.
##########
datafusion/optimizer/src/decorrelate_predicate_subquery.rs:
##########
@@ -613,34 +666,54 @@ fn build_join(
// Projecting the constant as a column of the outer side turns the
predicate
// into a real equi-join key, which fixes all three. It only pays for
itself
// on a join that ends up null-aware.
+ //
+ // Every constant element of a tuple is projected the same way, so that all
+ // value keys lead the equi-join keys in order.
let mut projected_left = None;
- if let Some((value, right_col, mut value_name, correlation_opt)) =
in_value_expr
- && value.column_refs().is_empty()
- && null_aware
+ if null_aware
+ && in_values
+ .iter()
+ .any(|(value, _)| value.column_refs().is_empty())
{
- // The projected column is unqualified, so a left field that already
has
- // this name — however unlikely — would make the reference ambiguous.
let left_schema = left.schema();
- while left_schema.fields().iter().any(|f| f.name() == &value_name) {
- value_name.push('_');
- }
- let value_col = Column::new_unqualified(value_name);
- let projections = left_schema
+ let mut projections = left_schema
.columns()
.into_iter()
.map(Expr::from)
- .chain(std::iter::once(value.alias(value_col.name())))
.collect::<Vec<_>>();
+ let is_tuple = in_values.len() > 1;
+ for (i, (value, _)) in in_values.iter_mut().enumerate() {
+ if !value.column_refs().is_empty() {
+ continue;
+ }
+ let mut value_name = if is_tuple {
+ format!("{alias}_value_{i}")
+ } else {
+ format!("{alias}_value")
+ };
+ // The projected column is unqualified, so a left field that
already
+ // has this name — however unlikely — would make the reference
+ // ambiguous.
+ while left_schema.fields().iter().any(|f| f.name() == &value_name)
{
+ value_name.push('_');
+ }
+ let value_col = Column::new_unqualified(value_name);
+ let constant = std::mem::replace(value,
Expr::Column(value_col.clone()));
+ projections.push(constant.alias(value_col.name()));
+ }
projected_left = Some(
LogicalPlanBuilder::from(left.clone())
.project(projections)?
.build()?,
);
- // Rebuild the `IN` equality against the projected column.
- let in_predicate = Expr::eq(Expr::Column(value_col),
Expr::Column(right_col));
- join_filter = in_predicate_first(in_predicate, correlation_opt);
+ // Rebuild the `IN` equalities against the projected columns.
+ join_filter = in_predicate_first(in_values_predicate(&in_values),
in_correlation);
}
let left = projected_left.as_ref().unwrap_or(left);
+ // The `IN` equalities lead the join filter, so after they are extracted
+ // into equi-join keys the first `null_aware_value_keys` keys are the
+ // `NOT IN` value keys (see `Join::null_aware_value_keys`).
+ let null_aware_value_keys = in_values.len().max(1);
Review Comment:
Agreed. `build_join` now checks, before building a null-aware join, that
every `NOT IN` value pair passes `find_valid_equijoin_key_pair` and `can_hash`
on both sides, and returns `not_impl_err!` otherwise. Since the check runs on
the value pairs regardless of their count, it also closes the scalar gap you
found on `main`. Both the tuple and the scalar RunEndEncoded queries are
covered in `null_aware_multi_column.slt`.
##########
datafusion/expr/src/expr.rs:
##########
@@ -1399,6 +1399,34 @@ impl InSubquery {
negated,
}
}
+
+ /// The tuple elements of a multi-column `(a, b, ...) IN (SELECT x, y,
...)`,
+ /// or `None` for a single-column `IN`. See [`in_subquery_tuple_values`].
+ pub fn tuple_values(&self) -> Option<&[Expr]> {
+ in_subquery_tuple_values(&self.expr, &self.subquery.subquery)
+ }
+}
+
+/// The tuple elements of a multi-column `(a, b, ...) IN (SELECT x, y, ...)`,
+/// given the compared expression `expr` and the `subquery` plan, or `None`
+/// for a single-column `IN`.
+///
+/// The tuple is planned as a `struct` call. It is a multi-column `IN` only
when
+/// the subquery returns more than one column; against a single column the
+/// `struct` is one value compared with a struct-typed column. The number of
+/// elements is not checked against the number of subquery columns here.
+pub fn in_subquery_tuple_values<'a>(
+ expr: &'a Expr,
+ subquery: &crate::LogicalPlan,
+) -> Option<&'a [Expr]> {
+ match expr {
+ Expr::ScalarFunction(func)
+ if func.func.name() == "struct" &&
subquery.schema().fields().len() > 1 =>
Review Comment:
Fixed. `in_subquery_tuple_values` now also recognizes a
`ScalarValue::Struct` literal and returns its elements as owned exprs
(`Cow<[Expr]>`). `(1, 2) NOT IN (...)` is covered in
`null_aware_multi_column.slt`, both as a filter and as projected values
(matched / UNKNOWN / mismatched).
##########
datafusion/physical-plan/src/joins/hash_join/exec.rs:
##########
@@ -245,8 +256,18 @@ impl NullAwareMode {
join_type: JoinType,
partition_mode: PartitionMode,
num_keys: usize,
+ num_value_keys: usize,
has_filter: bool,
) -> Result<Self> {
+ if num_value_keys == 0 || num_value_keys > num_keys {
+ return plan_err!(
+ "null_aware {join_type} join needs between 1 and {num_keys}
`NOT IN` value keys, got {num_value_keys}"
+ );
+ }
+ // `num_keys > 1` also covers a multi-column value key without any
+ // correlation: a NULL in one tuple element leaves the comparison FALSE
+ // whenever another element is a definite mismatch, so it too must be
+ // decided per build row.
let correlated = num_keys > 1 || has_filter;
Review Comment:
Thanks for the measurements. I'd like to keep the build-side change and the
narrowing as a follow-up, as you suggested. For now:
- The cost is documented on `HashJoinExec::null_aware_value_keys`: the outer
side stays the build side, since only a single-key join is swapped, and without
correlation keys every NULL-valued row is paired with every row on the other
side.
- The `null_aware_join` suite has three new queries for these shapes: Q10
with no NULLs, Q11 with 1% NULLs in one subquery element, and Q12 with 1% NULLs
in one outer element.
Treating an element that is NOT NULL on both sides as a scope key, and a
multi-key null-aware `RightAnti` that builds on the subquery side: agreed these
are the right next steps; leaving them as follow-ups.
--
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]