szehon-ho commented on code in PR #58948:
URL: https://github.com/apache/spark/pull/58948#discussion_r4078430910
##########
sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/PushDownUtils.scala:
##########
@@ -144,6 +156,28 @@ object PushDownUtils extends Logging {
}
}
+ // Extract a necessary condition from a deterministic filter that cannot be
fully translated.
+ // The caller must retain the original filter for post-scan evaluation. AND
can use either
+ // child, but OR requires both children. Other expressions, including NOT,
must translate whole.
+ private def extractPushablePredicate(
+ expression: Expression,
+ canTranslate: Expression => Boolean): Option[Expression] = expression
match {
+ case And(left, right) =>
+ val l = extractPushablePredicate(left, canTranslate)
+ val r = extractPushablePredicate(right, canTranslate)
+ (l, r) match {
+ case (Some(a), Some(b)) => Some(And(a, b))
+ case _ => l.orElse(r)
+ }
+ case Or(left, right) =>
+ for {
+ l <- extractPushablePredicate(left, canTranslate)
+ r <- extractPushablePredicate(right, canTranslate)
+ } yield Or(l, r)
+ case other if canTranslate(other) => Some(other)
+ case _ => None
+ }
Review Comment:
Could we accept an `Option`-returning translator directly?
```suggestion
private def extractPushablePredicate[T](
expression: Expression,
translate: Expression => Option[T]): Option[Expression] = expression
match {
case And(left, right) =>
val l = extractPushablePredicate(left, translate)
val r = extractPushablePredicate(right, translate)
(l, r) match {
case (Some(a), Some(b)) => Some(And(a, b))
case _ => l.orElse(r)
}
case Or(left, right) =>
for {
l <- extractPushablePredicate(left, translate)
r <- extractPushablePredicate(right, translate)
} yield Or(l, r)
case other =>
translate(other).map(_ => other)
}
```
This lets callers pass translation functions directly and handles their
`Option` results with `map`, without using `isDefined` in the callbacks.
##########
sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/PushDownUtils.scala:
##########
@@ -105,12 +112,17 @@ object PushDownUtils extends Logging {
val translatedFilters = mutable.ArrayBuffer.empty[Predicate]
val untranslatableExprs = mutable.ArrayBuffer.empty[Expression]
+ def translateFilter(expression: Expression): Option[Predicate] = {
+ DataSourceV2Strategy.translateFilterV2WithMapping(
+ expression, Some(translatedFilterToExpr))
+ }
+
for (filterExpr <- deterministicFilters) {
- val translated =
- DataSourceV2Strategy.translateFilterV2WithMapping(
- filterExpr, Some(translatedFilterToExpr))
+ val translated = translateFilter(filterExpr)
if (translated.isEmpty) {
untranslatableExprs += filterExpr
+ extractPushablePredicate(filterExpr, e =>
translateFilter(e).isDefined)
+ .flatMap(translateFilter).foreach(translatedFilters += _)
} else {
translatedFilters += translated.get
}
Review Comment:
With the updated helper, could we pattern-match on the translation result
here?
```suggestion
translateFilter(filterExpr) match {
case Some(filter) =>
translatedFilters += filter
case None =>
untranslatableExprs += filterExpr
extractPushablePredicate(
filterExpr,
DataSourceV2Strategy.translateFilterV2)
.flatMap(translateFilter)
.foreach(translatedFilters += _)
}
```
- Makes the success and fallback paths explicit, avoiding `isEmpty` followed
by `.get`.
- Uses translation without mapping updates when checking extraction
candidates.
##########
sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/PushDownUtils.scala:
##########
@@ -73,12 +73,19 @@ object PushDownUtils extends Logging {
// Catalyst filter expression that can't be translated to data source
filters.
val untranslatableExprs = mutable.ArrayBuffer.empty[Expression]
+ def translateFilter(expression: Expression): Option[sources.Filter] = {
+ DataSourceStrategy.translateFilterWithMapping(expression,
Some(translatedFilterToExpr),
+ nestedPredicatePushdownEnabled = true)
+ }
+
for (filterExpr <- filters) {
- val translated =
- DataSourceStrategy.translateFilterWithMapping(filterExpr,
Some(translatedFilterToExpr),
- nestedPredicatePushdownEnabled = true)
+ val translated = translateFilter(filterExpr)
if (translated.isEmpty) {
untranslatableExprs += filterExpr
+ if (filterExpr.deterministic) {
+ extractPushablePredicate(filterExpr, e =>
translateFilter(e).isDefined)
+ .flatMap(translateFilter).foreach(translatedFilters += _)
+ }
} else {
translatedFilters += translated.get
}
Review Comment:
The corresponding V1 change, keeping the determinism guard and nested-column
support:
```suggestion
translateFilter(filterExpr) match {
case Some(filter) =>
translatedFilters += filter
case None =>
untranslatableExprs += filterExpr
if (filterExpr.deterministic) {
extractPushablePredicate(
filterExpr,
e => DataSourceStrategy.translateFilter(
e, supportNestedPredicatePushdown = true))
.flatMap(translateFilter)
.foreach(translatedFilters += _)
}
}
```
--
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]