This is an automated email from the ASF dual-hosted git repository.

codope pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/hudi.git


The following commit(s) were added to refs/heads/master by this push:
     new 22ae3bad95a [HUDI-8598] Fix AND operator to not filter unsupported 
query types (#12363)
22ae3bad95a is described below

commit 22ae3bad95a6528f43629af2dccf037248603857
Author: Lokesh Jain <[email protected]>
AuthorDate: Fri Nov 29 08:52:38 2024 +0530

    [HUDI-8598] Fix AND operator to not filter unsupported query types (#12363)
---
 .../org/apache/hudi/RecordLevelIndexSupport.scala  | 45 ++++++++++------------
 .../org/apache/hudi/SparkBaseIndexSupport.scala    |  8 +---
 .../apache/hudi/TestRecordLevelIndexSupport.scala  |  5 +++
 3 files changed, 27 insertions(+), 31 deletions(-)

diff --git 
a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/RecordLevelIndexSupport.scala
 
b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/RecordLevelIndexSupport.scala
index 818228bcf6d..454014b9728 100644
--- 
a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/RecordLevelIndexSupport.scala
+++ 
b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/RecordLevelIndexSupport.scala
@@ -25,8 +25,8 @@ import org.apache.hudi.common.model.FileSlice
 import org.apache.hudi.common.model.HoodieRecord.HoodieMetadataField
 import org.apache.hudi.common.model.HoodieTableQueryType.SNAPSHOT
 import org.apache.hudi.common.table.HoodieTableMetaClient
-import 
org.apache.hudi.common.table.timeline.InstantComparison.compareTimestamps
 import org.apache.hudi.common.table.timeline.InstantComparison
+import 
org.apache.hudi.common.table.timeline.InstantComparison.compareTimestamps
 import org.apache.hudi.keygen.KeyGenUtils
 import org.apache.hudi.metadata.HoodieTableMetadataUtil
 import org.apache.hudi.storage.StoragePath
@@ -148,19 +148,15 @@ object RecordLevelIndexSupport {
 
   def filterQueryWithRecordKey(queryFilter: Expression, recordKeyOpt: 
Option[String],
                                literalGenerator: Function2[AttributeReference, 
Literal, String]): Option[(Expression, List[String])] = {
-    filterQueryWithRecordKey(queryFilter, recordKeyOpt, literalGenerator, 
getDefaultAttributeFetcher())
+    filterQueryWithRecordKey(queryFilter, recordKeyOpt, literalGenerator, 
getDefaultAttributeFetcher())._1
   }
 
   def filterQueryWithRecordKey(queryFilter: Expression, recordKeyOpt: 
Option[String], attributeFetcher: Function1[Expression, Expression]): 
Option[(Expression, List[String])] = {
-    filterQueryWithRecordKey(queryFilter, recordKeyOpt, 
getSimpleLiteralGenerator(), attributeFetcher)
-  }
-
-  def isSupported(expr: Expression): Boolean = {
-    expr.isInstanceOf[In] || expr.isInstanceOf[EqualTo] || 
expr.isInstanceOf[And]
+    filterQueryWithRecordKey(queryFilter, recordKeyOpt, 
getSimpleLiteralGenerator(), attributeFetcher)._1
   }
 
   def filterQueryWithRecordKey(queryFilter: Expression, recordKeyOpt: 
Option[String], literalGenerator: Function2[AttributeReference, Literal, 
String],
-                               attributeFetcher: Function1[Expression, 
Expression]): Option[(Expression, List[String])] = {
+                               attributeFetcher: Function1[Expression, 
Expression]): (Option[(Expression, List[String])], Boolean) = {
     queryFilter match {
       case equalToQuery: EqualTo =>
         val attributeLiteralTuple = 
getAttributeLiteralTuple(attributeFetcher.apply(equalToQuery.left), 
attributeFetcher.apply(equalToQuery.right)).orNull
@@ -169,12 +165,12 @@ object RecordLevelIndexSupport {
           val literal = attributeLiteralTuple._2
           if (attribute != null && attribute.name != null && 
attributeMatchesRecordKey(attribute.name, recordKeyOpt)) {
             val recordKeyLiteral = literalGenerator.apply(attribute, literal)
-            Option.apply(equalToQuery, List.apply(recordKeyLiteral))
+            (Option.apply(equalToQuery, List.apply(recordKeyLiteral)), true)
           } else {
-            Option.empty
+            (Option.empty, true)
           }
         } else {
-          Option.empty
+          (Option.empty, true)
         }
 
       case inQuery: In =>
@@ -200,35 +196,36 @@ object RecordLevelIndexSupport {
           case _ => validINQuery = false
         }
         if (validINQuery) {
-          Option.apply(inQuery, literals)
+          (Option.apply(inQuery, literals), true)
         } else {
-          Option.empty
+          (Option.empty, true)
         }
 
       // Handle And expression (composite filter)
       case andQuery: And =>
-        if (!isSupported(andQuery.left) || !isSupported(andQuery.right)) {
-          Option.empty
-        } else {
-          val leftResult = filterQueryWithRecordKey(andQuery.left, 
recordKeyOpt, literalGenerator, attributeFetcher)
-          val rightResult = filterQueryWithRecordKey(andQuery.right, 
recordKeyOpt, literalGenerator, attributeFetcher)
+        val leftResult = filterQueryWithRecordKey(andQuery.left, recordKeyOpt, 
literalGenerator, attributeFetcher)
+        val rightResult = filterQueryWithRecordKey(andQuery.right, 
recordKeyOpt, literalGenerator, attributeFetcher)
 
+        val isSupported = leftResult._2 && rightResult._2
+        if (!isSupported) {
+          (Option.empty, false)
+        } else {
           // If both left and right filters are valid, concatenate their 
results
-          (leftResult, rightResult) match {
+          (leftResult._1, rightResult._1) match {
             case (Some((leftExp, leftKeys)), Some((rightExp, rightKeys))) =>
               // Return concatenated expressions and record keys
-              Option.apply(And(leftExp, rightExp), leftKeys ++ rightKeys)
+              (Option.apply(And(leftExp, rightExp), leftKeys ++ rightKeys), 
true)
             case (Some((leftExp, leftKeys)), None) =>
               // Return concatenated expressions and record keys
-              Option.apply(leftExp, leftKeys)
+              (Option.apply(leftExp, leftKeys), true)
             case (None, Some((rightExp, rightKeys))) =>
               // Return concatenated expressions and record keys
-              Option.apply(rightExp, rightKeys)
-            case _ => Option.empty
+              (Option.apply(rightExp, rightKeys), true)
+            case _ => (Option.empty, true)
           }
         }
 
-      case _ => Option.empty
+      case _ => (Option.empty, false)
     }
   }
 
diff --git 
a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/SparkBaseIndexSupport.scala
 
b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/SparkBaseIndexSupport.scala
index fed55c29eff..41eb0978e70 100644
--- 
a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/SparkBaseIndexSupport.scala
+++ 
b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/SparkBaseIndexSupport.scala
@@ -177,13 +177,7 @@ abstract class SparkBaseIndexSupport(spark: SparkSession,
                   }
                 ).foreach {
                   case (exp: Expression, recKeys: List[String]) =>
-                    exp match {
-                      // For IN, add each element individually to recordKeys
-                      case _: In => recordKeys = recordKeys ++ recKeys
-
-                      // For other cases, basically EqualTo, concatenate 
recKeys with the default separator
-                      case _ => recordKeys = recordKeys ++ recKeys
-                    }
+                    recordKeys = recordKeys ++ recKeys
                     recordKeyQueries = recordKeyQueries :+ exp
                 }
               }
diff --git 
a/hudi-spark-datasource/hudi-spark-common/src/test/scala/org/apache/hudi/TestRecordLevelIndexSupport.scala
 
b/hudi-spark-datasource/hudi-spark-common/src/test/scala/org/apache/hudi/TestRecordLevelIndexSupport.scala
index 7ff7e9488ef..9ee8ba7d940 100644
--- 
a/hudi-spark-datasource/hudi-spark-common/src/test/scala/org/apache/hudi/TestRecordLevelIndexSupport.scala
+++ 
b/hudi-spark-datasource/hudi-spark-common/src/test/scala/org/apache/hudi/TestRecordLevelIndexSupport.scala
@@ -125,5 +125,10 @@ class TestRecordLevelIndexSupport {
     testFilter = And(rk3InFilter, Or(rk2EqFilter, rk1EqFilter))
     result = RecordLevelIndexSupport.filterQueryWithRecordKey(testFilter, 
Option.apply("rk1"), RecordLevelIndexSupport.getComplexKeyLiteralGenerator())
     assertTrue(result.isEmpty)
+
+    // Case 11: Test unsupported query type with And. Here the unsupported 
query OR is at second level from the root query type
+    testFilter = And(rk1EqFilter, And(rk1InFilter, Or(rk1EqFilter, 
rk1EqFilter)))
+    result = RecordLevelIndexSupport.filterQueryWithRecordKey(testFilter, 
Option.apply("rk1"), RecordLevelIndexSupport.getComplexKeyLiteralGenerator())
+    assertTrue(result.isEmpty)
   }
 }

Reply via email to